From 6ec5b76c54a500748791c5832865aea731569139 Mon Sep 17 00:00:00 2001 From: mptyl Date: Tue, 18 Aug 2026 02:43:42 +0200 Subject: [PATCH] fix(auth): serialize restore checkpoint lifecycle --- tools/tht/internal/backup/create.go | 14 +- tools/tht/internal/backup/restore.go | 75 +++-- tools/tht/internal/backup/restore_host.go | 4 +- tools/tht/internal/backup/restore_test.go | 338 +++++++++++++++++++++- 4 files changed, 391 insertions(+), 40 deletions(-) diff --git a/tools/tht/internal/backup/create.go b/tools/tht/internal/backup/create.go index 400d2d0a..73999343 100644 --- a/tools/tht/internal/backup/create.go +++ b/tools/tht/internal/backup/create.go @@ -97,8 +97,12 @@ type dependencies struct { // Create creates an archive with the real Docker command boundary. It performs no shell // interpolation and never supplies secret contents in process arguments. func Create(ctx context.Context, installation config.Installation, request CreateRequest) (Result, error) { + return createWithDependencies(ctx, installation, request, productionCreateDependencies(installation)) +} + +func productionCreateDependencies(installation config.Installation) dependencies { runner := hostRunner{runner: compose.NewRunner(""), binary: "docker", profile: installation.Profile} - return createWithDependencies(ctx, installation, request, dependencies{ + return dependencies{ runner: runner, now: time.Now, homeDir: os.UserHomeDir, @@ -113,7 +117,7 @@ func Create(ctx context.Context, installation config.Installation, request Creat sleep: time.Sleep, reserveOutput: reserveArchiveOutput, publishReserved: publishReservedArchive, - }) + } } func createWithDependencies(ctx context.Context, installation config.Installation, request CreateRequest, dependencies dependencies) (result Result, resultErr error) { @@ -133,7 +137,13 @@ func createWithDependencies(ctx context.Context, installation config.Installatio resultErr = errors.Join(resultErr, fmt.Errorf("release backup lifecycle lock: %w", releaseErr)) } }() + return createWithDependenciesLockHeld(ctx, installation, request, dependencies) +} +// createWithDependenciesLockHeld performs backup creation while the caller owns the installation +// lifecycle lock. It must never acquire a lifecycle lock itself: Restore uses this primitive to +// create its recovery checkpoint inside its already-locked transaction. +func createWithDependenciesLockHeld(ctx context.Context, installation config.Installation, request CreateRequest, dependencies dependencies) (result Result, resultErr error) { if err := rejectInlineSecretValues(installation.EnvFile); err != nil { return Result{}, err } diff --git a/tools/tht/internal/backup/restore.go b/tools/tht/internal/backup/restore.go index 8c9aa4d2..e90022dc 100644 --- a/tools/tht/internal/backup/restore.go +++ b/tools/tht/internal/backup/restore.go @@ -33,8 +33,11 @@ type restoreLock interface{ Release() error } type restoreVerify func(context.Context, config.Installation, archiveRunner) error type restoreDependencies struct { - preflight func(context.Context, config.Installation, PreflightRequest) (PreflightResult, error) - checkpoint func(context.Context, config.Installation, CreateRequest) (Result, error) + preflight func(context.Context, config.Installation, PreflightRequest) (PreflightResult, error) + // checkpointLocked creates the secret-aware recovery archive while the caller already owns + // the installation lifecycle lock. It must not call public Create, which would re-acquire the + // non-reentrant lock and deadlock the restore transaction. + checkpointLocked func(context.Context, config.Installation, CreateRequest) (Result, error) prepareRecovery func(context.Context, config.Installation, string) (PreflightResult, error) recover func(context.Context, config.Installation, PreflightResult, bool) error cleanupCheckpoint func(string) error @@ -81,36 +84,25 @@ func restoreWithDependencies(ctx context.Context, installation config.Installati if request.Archive == "" { return RestoreResult{}, errors.New("restore archive is required") } - if deps.preflight == nil || deps.checkpoint == nil || deps.prepareRecovery == nil || deps.recover == nil || deps.cleanupCheckpoint == nil || deps.acquireLock == nil || deps.runner == nil || deps.restoreFile == nil || deps.restoreVolume == nil || deps.resetAuthenticationState == nil || deps.verify == nil { + if deps.preflight == nil || deps.checkpointLocked == nil || deps.prepareRecovery == nil || deps.recover == nil || deps.cleanupCheckpoint == nil || deps.acquireLock == nil || deps.runner == nil || deps.restoreFile == nil || deps.restoreVolume == nil || deps.resetAuthenticationState == nil || deps.verify == nil { return RestoreResult{}, errors.New("restore dependencies are incomplete") } preflight, err := deps.preflight(ctx, installation, PreflightRequest{Archive: request.Archive, Confirm: true, AllowExternalSecrets: true}) if err != nil { return RestoreResult{}, err } - defer preflight.CloseArchive() archive, err := preflight.RevalidateArchive() if err != nil { + _ = preflight.CloseArchive() return RestoreResult{}, err } - checkpoint, err := deps.checkpoint(ctx, installation, CreateRequest{IncludeSecrets: true, Confirm: true}) - if err != nil { - return RestoreResult{}, fmt.Errorf("create recovery checkpoint: %w", err) - } - recovery, err := deps.prepareRecovery(ctx, installation, checkpoint.Path) - if err != nil { - cleanupErr := deps.cleanupCheckpoint(checkpoint.Path) - return RestoreResult{}, errors.Join(fmt.Errorf("validate recovery checkpoint: %w", err), cleanupErr) - } - defer func() { - if cleanupErr := deps.cleanupCheckpoint(checkpoint.Path); cleanupErr != nil { - resultErr = errors.Join(resultErr, fmt.Errorf("destroy recovery checkpoint: %w", cleanupErr)) - } - }() - defer recovery.CloseArchive() + // Lock ordering is lifecycle lock -> Compose/operator maintenance barrier. The core operator + // command never acquires the host lifecycle lock, so this order cannot form a lock cycle with + // Docker Compose or the durable maintenance marker. lock, err := deps.acquireLock(installation) if err != nil { + _ = preflight.CloseArchive() return result, err } defer func() { @@ -119,12 +111,20 @@ func restoreWithDependencies(ctx context.Context, installation config.Installati resultErr = errors.Join(resultErr, fmt.Errorf("release restore lifecycle lock: %w", releaseErr)) } }() + defer preflight.CloseArchive() - wasRunning, err := installationRunning(ctx, installation, deps.runner) + checkpoint, err := deps.checkpointLocked(ctx, installation, CreateRequest{IncludeSecrets: true, Confirm: true}) if err != nil { - return result, err + return result, fmt.Errorf("create recovery checkpoint: %w", err) } - state := restoreTransactionState{wasRunning: wasRunning} + recovery, err := deps.prepareRecovery(ctx, installation, checkpoint.Path) + if err != nil { + cleanupErr := deps.cleanupCheckpoint(checkpoint.Path) + return result, errors.Join(fmt.Errorf("validate recovery checkpoint: %w", err), cleanupErr) + } + defer recovery.CloseArchive() + + state := restoreTransactionState{} defer func() { if state.recoveryRequired(resultErr) { recoveryContext, cancel := boundedCleanupContext() @@ -141,10 +141,27 @@ func restoreWithDependencies(ctx context.Context, installation config.Installati state.stopAttempted = false } } + + checkpointCleanupSucceeded := true + var cleanupErr error + if checkpointErr := deps.cleanupCheckpoint(checkpoint.Path); checkpointErr != nil { + checkpointCleanupSucceeded = false + cleanupErr = errors.Join(cleanupErr, fmt.Errorf("destroy recovery checkpoint: %w", checkpointErr)) + } if !state.maintenanceAttempted { + if cleanupErr != nil { + result = RestoreResult{} + resultErr = errors.Join(resultErr, cleanupErr) + } + return + } + // A recovery checkpoint remains secret-bearing transaction state. Do not reopen + // admissions until it has been safely deleted, even if the restored target verified. + if !checkpointCleanupSucceeded { + result = RestoreResult{} + resultErr = errors.Join(resultErr, cleanupErr) return } - var cleanupErr error restartCompleted := true if state.wasRunning && state.stopAttempted && state.mayDeactivateMaintenance() { startErr, started := retryBoundedCleanup(func(cleanupContext context.Context) error { @@ -177,6 +194,12 @@ func restoreWithDependencies(ctx context.Context, installation config.Installati resultErr = errors.Join(resultErr, cleanupErr) } }() + + wasRunning, err := installationRunning(ctx, installation, deps.runner) + if err != nil { + return result, err + } + state.wasRunning = wasRunning if state.wasRunning { // The activation command may take effect even when its response is lost. Track the attempt, // not merely a successful return, so every subsequent path compensates from durable state. @@ -213,12 +236,6 @@ func restoreWithDependencies(ctx context.Context, installation config.Installati } state.verified = true result.Verified = true - if state.wasRunning { - if err := maintenance(ctx, installation, deps.runner, false); err != nil { - return result, err - } - state.maintenanceAttempted = false - } return result, nil } diff --git a/tools/tht/internal/backup/restore_host.go b/tools/tht/internal/backup/restore_host.go index 8519886d..e55f364b 100644 --- a/tools/tht/internal/backup/restore_host.go +++ b/tools/tht/internal/backup/restore_host.go @@ -41,13 +41,13 @@ func productionRestoreDependencies(installation config.Installation) restoreDepe }, }) }, - checkpoint: func(ctx context.Context, target config.Installation, request CreateRequest) (Result, error) { + checkpointLocked: func(ctx context.Context, target config.Installation, request CreateRequest) (Result, error) { path, err := restoreCheckpointPath(target, time.Now().UTC()) if err != nil { return Result{}, err } request.Output = path - return Create(ctx, target, request) + return createWithDependenciesLockHeld(ctx, target, request, productionCreateDependencies(target)) }, cleanupCheckpoint: cleanupRecoveryCheckpoint, acquireLock: func(target config.Installation) (restoreLock, error) { diff --git a/tools/tht/internal/backup/restore_test.go b/tools/tht/internal/backup/restore_test.go index 7181c8d4..ddba356d 100644 --- a/tools/tht/internal/backup/restore_test.go +++ b/tools/tht/internal/backup/restore_test.go @@ -15,6 +15,7 @@ import ( "github.com/aritmolab/thothii/tools/tht/internal/compose" "github.com/aritmolab/thothii/tools/tht/internal/config" + "github.com/aritmolab/thothii/tools/tht/internal/lifecycle" ) func TestRestorePublicPathUsesConcreteProductionPreflight(t *testing.T) { @@ -119,7 +120,7 @@ func TestRestoreStoppedInstallationRunsCheckpointRestoreAndVerification(t *testi var events []string var checkpointRequest CreateRequest deps := restoreTestDependencies(t, runner) - deps.checkpoint = func(_ context.Context, _ config.Installation, request CreateRequest) (Result, error) { + deps.checkpointLocked = func(_ context.Context, _ config.Installation, request CreateRequest) (Result, error) { events = append(events, "checkpoint") checkpointRequest = request return Result{Path: "/tmp/checkpoint.zip"}, nil @@ -158,11 +159,314 @@ func TestRestoreStoppedInstallationRunsCheckpointRestoreAndVerification(t *testi if !checkpointRequest.IncludeSecrets || !checkpointRequest.Confirm { t.Fatalf("checkpoint request = %#v, want private confirmed secret-aware checkpoint", checkpointRequest) } - if got, want := events, []string{"checkpoint", "prepare-recovery", "lock", "file:configuration/operator.env", "health", "doctor", "pi", "workspace", "unlock", "cleanup-checkpoint"}; !equalStrings(got, want) { + if got, want := events, []string{"lock", "checkpoint", "prepare-recovery", "file:configuration/operator.env", "health", "doctor", "pi", "workspace", "cleanup-checkpoint", "unlock"}; !equalStrings(got, want) { t.Fatalf("restore events = %v, want %v", got, want) } } +func TestRestoreClosesTargetArchiveBeforeReleasingLifecycleLock(t *testing.T) { + installation := preflightTestInstallation(t) + archive := filepath.Join(t.TempDir(), "restore.zip") + writePreflightArchive(t, archive, preflightArchiveSpec{ + entries: []preflightArchiveEntry{{path: "configuration/operator.env", body: []byte("safe")}}, + }) + + runner := newBackupRunner(installation, false) + deps := restoreTestDependencies(t, runner) + var target PreflightResult + deps.preflight = func(ctx context.Context, targetInstallation config.Installation, request PreflightRequest) (PreflightResult, error) { + result, err := Preflight(ctx, targetInstallation, request, permissivePreflightDependencies()) + target = result + return result, err + } + unlockBeforeTargetClose := false + deps.acquireLock = func(config.Installation) (restoreLock, error) { + return fakeRestoreLock{release: func() { + if target.archive != nil && target.archive.file != nil { + if _, err := target.archive.file.Stat(); err == nil { + unlockBeforeTargetClose = true + } + } + }}, nil + } + + if _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: archive, Confirm: true}, deps); err != nil { + t.Fatal(err) + } + if unlockBeforeTargetClose { + t.Fatal("restore released its lifecycle lock before closing the target archive") + } +} + +func TestRestoreLifecycleLockExcludesCompetingTransactionsUntilTerminalCleanup(t *testing.T) { + targetFailure := errors.New("target restore failed") + recoveryFailure := errors.New("recovery restore failed") + for _, scenario := range []struct { + name string + stages []string + configure func(*restoreDependencies, *lifecycleGateRunner, func(string), context.CancelFunc) + wantErrors []error + wantVerified bool + wantMaintenance bool + }{ + { + name: "successful target", + stages: []string{"checkpoint", "target", "verification", "checkpoint-cleanup", "final-barrier-release"}, + wantVerified: true, + wantMaintenance: false, + configure: func(deps *restoreDependencies, _ *lifecycleGateRunner, gate func(string), _ context.CancelFunc) { + deps.restoreFile = func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error { + gate("target") + return nil + } + deps.verify["health"] = func(context.Context, config.Installation, archiveRunner) error { + gate("verification") + return nil + } + }, + }, + { + name: "target failure with verified recovery", + stages: []string{"checkpoint", "target-failure", "recovery", "checkpoint-cleanup", "final-barrier-release"}, + wantErrors: []error{targetFailure}, + wantMaintenance: false, + configure: func(deps *restoreDependencies, runner *lifecycleGateRunner, gate func(string), _ context.CancelFunc) { + deps.restoreFile = func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error { + gate("target-failure") + return targetFailure + } + deps.recover = func(context.Context, config.Installation, PreflightResult, bool) error { + gate("recovery") + runner.running, runner.coreRunning = true, true + return nil + } + }, + }, + { + name: "recovery failure", + stages: []string{"checkpoint", "target-failure", "recovery-failure", "checkpoint-cleanup"}, + wantErrors: []error{targetFailure, recoveryFailure}, + wantMaintenance: true, + configure: func(deps *restoreDependencies, _ *lifecycleGateRunner, gate func(string), _ context.CancelFunc) { + deps.restoreFile = func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error { + gate("target-failure") + return targetFailure + } + deps.recover = func(context.Context, config.Installation, PreflightResult, bool) error { + gate("recovery-failure") + return recoveryFailure + } + }, + }, + { + name: "caller cancellation", + stages: []string{"checkpoint", "target-cancel", "recovery", "checkpoint-cleanup", "final-barrier-release"}, + wantErrors: []error{context.Canceled}, + wantMaintenance: false, + configure: func(deps *restoreDependencies, runner *lifecycleGateRunner, gate func(string), cancel context.CancelFunc) { + deps.restoreFile = func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error { + gate("target-cancel") + cancel() + return context.Canceled + } + deps.recover = func(context.Context, config.Installation, PreflightResult, bool) error { + gate("recovery") + runner.running, runner.coreRunning = true, true + return nil + } + }, + }, + } { + scenario := scenario + t.Run(scenario.name, func(t *testing.T) { + fixture := newBackupFixture(t, "local") + installation := fixture.installation + archive := filepath.Join(t.TempDir(), "restore.zip") + writePreflightArchive(t, archive, preflightArchiveSpec{ + installationID: fixture.installationID, + entries: []preflightArchiveEntry{{path: "configuration/operator.env", body: []byte("safe")}}, + }) + runner := &lifecycleGateRunner{fakeBackupRunner: newBackupRunner(installation, true)} + stages := make(chan string) + continueStage := make(chan struct{}) + gate := func(stage string) { + stages <- stage + <-continueStage + } + runner.beforeFinalMaintenanceRelease = func() { gate("final-barrier-release") } + caller, cancel := context.WithCancel(context.Background()) + defer cancel() + deps := restoreTestDependencies(t, runner) + deps.acquireLock = func(target config.Installation) (restoreLock, error) { return lifecycle.Acquire(target) } + deps.checkpointLocked = func(context.Context, config.Installation, CreateRequest) (Result, error) { + gate("checkpoint") + return Result{Path: filepath.Join(t.TempDir(), "checkpoint.zip")}, nil + } + deps.cleanupCheckpoint = func(string) error { + gate("checkpoint-cleanup") + return nil + } + scenario.configure(&deps, runner, gate, cancel) + + type outcome struct { + result RestoreResult + err error + } + done := make(chan outcome, 1) + go func() { + result, err := restoreWithDependencies(caller, installation, RestoreRequest{Archive: archive, Confirm: true}, deps) + done <- outcome{result: result, err: err} + }() + + var failures []error + for _, wantStage := range scenario.stages { + select { + case stage := <-stages: + if stage != wantStage { + failures = append(failures, fmt.Errorf("lifecycle stage = %q, want %q", stage, wantStage)) + } + case <-time.After(2 * time.Second): + failures = append(failures, fmt.Errorf("timed out waiting for lifecycle stage %q", wantStage)) + } + if err := competingRestoreAndBackupEntry(installation, archive, t); err != nil { + failures = append(failures, fmt.Errorf("%s: %w", wantStage, err)) + } + continueStage <- struct{}{} + } + result := <-done + for _, wantErr := range scenario.wantErrors { + if !errors.Is(result.err, wantErr) { + failures = append(failures, fmt.Errorf("restore error = %v, want %v", result.err, wantErr)) + } + } + if result.result.Verified != scenario.wantVerified { + failures = append(failures, fmt.Errorf("restore verified = %t, want %t", result.result.Verified, scenario.wantVerified)) + } + if runner.maintenance != scenario.wantMaintenance { + failures = append(failures, fmt.Errorf("maintenance = %t, want %t", runner.maintenance, scenario.wantMaintenance)) + } + if err := lifecycleLockFreeAfterTerminalRestore(installation); err != nil { + failures = append(failures, err) + } + if len(failures) != 0 { + t.Fatal(errors.Join(failures...)) + } + }) + } +} + +func competingRestoreAndBackupEntry(installation config.Installation, archive string, t *testing.T) error { + t.Helper() + backupRunner := newBackupRunner(installation, false) + _, backupErr := createWithDependencies( + context.Background(), installation, + CreateRequest{Output: filepath.Join(t.TempDir(), "competing-backup.zip")}, + testDependencies(t, backupRunner), + ) + if !errors.Is(backupErr, lifecycle.ErrLocked) || len(backupRunner.calls) != 0 { + return fmt.Errorf("competing backup entered: error=%v docker-calls=%d", backupErr, len(backupRunner.calls)) + } + + checkpointCalled := false + restoreDeps := restoreTestDependencies(t, newBackupRunner(installation, false)) + restoreDeps.acquireLock = func(target config.Installation) (restoreLock, error) { return lifecycle.Acquire(target) } + restoreDeps.checkpointLocked = func(context.Context, config.Installation, CreateRequest) (Result, error) { + checkpointCalled = true + return Result{Path: filepath.Join(t.TempDir(), "competing-checkpoint.zip")}, nil + } + _, restoreErr := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: archive, Confirm: true}, restoreDeps) + if !errors.Is(restoreErr, lifecycle.ErrLocked) || checkpointCalled { + return fmt.Errorf("competing restore entered: error=%v checkpoint=%t", restoreErr, checkpointCalled) + } + return nil +} + +func lifecycleLockFreeAfterTerminalRestore(installation config.Installation) error { + lock, err := lifecycle.Acquire(installation) + if err != nil { + return fmt.Errorf("lifecycle lock leaked after terminal restore: %w", err) + } + if err := lock.Release(); err != nil { + return fmt.Errorf("release post-restore lifecycle lock: %w", err) + } + return nil +} + +func TestRestoreCannotApplyAStaleCheckpointOverAnInterleavedRestore(t *testing.T) { + fixture := newBackupFixture(t, "local") + installation := fixture.installation + archive := filepath.Join(t.TempDir(), "restore.zip") + writePreflightArchive(t, archive, preflightArchiveSpec{ + installationID: fixture.installationID, + entries: []preflightArchiveEntry{{path: "configuration/operator.env", body: []byte("safe")}}, + }) + + targetState := "before-checkpoint" + checkpointState := "" + recoveryObserved := "" + checkpointEntered := make(chan struct{}) + continueCheckpoint := make(chan struct{}) + firstRunner := newBackupRunner(installation, true) + firstDeps := restoreTestDependencies(t, firstRunner) + firstDeps.acquireLock = func(target config.Installation) (restoreLock, error) { return lifecycle.Acquire(target) } + firstDeps.checkpointLocked = func(context.Context, config.Installation, CreateRequest) (Result, error) { + checkpointState = targetState + close(checkpointEntered) + <-continueCheckpoint + return Result{Path: filepath.Join(t.TempDir(), "first-checkpoint.zip")}, nil + } + firstDeps.restoreFile = func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error { + return errors.New("first target mutation failed before changing state") + } + firstDeps.recover = func(context.Context, config.Installation, PreflightResult, bool) error { + recoveryObserved = targetState + targetState = checkpointState + firstRunner.running, firstRunner.coreRunning = true, true + return nil + } + + done := make(chan error, 1) + go func() { + _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: archive, Confirm: true}, firstDeps) + done <- err + }() + select { + case <-checkpointEntered: + case <-time.After(2 * time.Second): + t.Fatal("first restore did not begin its recovery checkpoint") + } + + interleavedCheckpoint := false + interleavedMutation := false + interleavedDeps := restoreTestDependencies(t, newBackupRunner(installation, false)) + interleavedDeps.acquireLock = func(target config.Installation) (restoreLock, error) { return lifecycle.Acquire(target) } + interleavedDeps.checkpointLocked = func(context.Context, config.Installation, CreateRequest) (Result, error) { + interleavedCheckpoint = true + return Result{Path: filepath.Join(t.TempDir(), "interleaved-checkpoint.zip")}, nil + } + interleavedDeps.restoreFile = func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error { + interleavedMutation = true + targetState = "later-restore-state" + return nil + } + _, interleavedErr := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: archive, Confirm: true}, interleavedDeps) + close(continueCheckpoint) + firstErr := <-done + + if !errors.Is(interleavedErr, lifecycle.ErrLocked) || interleavedCheckpoint || interleavedMutation { + t.Fatalf("interleaved restore was admitted: error=%v checkpoint=%t mutation=%t", interleavedErr, interleavedCheckpoint, interleavedMutation) + } + if firstErr == nil { + t.Fatal("first restore unexpectedly succeeded after its target mutation failed") + } + if recoveryObserved != "before-checkpoint" || targetState != "before-checkpoint" { + t.Fatalf("stale checkpoint recovery observed=%q final=%q; want no interleaved state to overwrite", recoveryObserved, targetState) + } + if err := lifecycleLockFreeAfterTerminalRestore(installation); err != nil { + t.Fatal(err) + } +} + func TestRestoreResetsAuthenticationStateBeforeRestart(t *testing.T) { installation := preflightTestInstallation(t) archive := restoreArchive(t) @@ -238,7 +542,7 @@ func TestRestorePreflightFailureDoesNotMutateTarget(t *testing.T) { return PreflightResult{}, preflightErr } checkpointCalls := 0 - deps.checkpoint = func(context.Context, config.Installation, CreateRequest) (Result, error) { + deps.checkpointLocked = func(context.Context, config.Installation, CreateRequest) (Result, error) { checkpointCalls++ return Result{}, nil } @@ -263,7 +567,7 @@ func TestRestoreCheckpointFailureDoesNotMutateTarget(t *testing.T) { runner := newBackupRunner(installation, true) deps := restoreTestDependencies(t, runner) checkpointErr := errors.New("checkpoint unavailable") - deps.checkpoint = func(context.Context, config.Installation, CreateRequest) (Result, error) { + deps.checkpointLocked = func(context.Context, config.Installation, CreateRequest) (Result, error) { return Result{}, checkpointErr } restoredFiles := 0 @@ -287,7 +591,7 @@ func TestRestoreFileFailureRollsBackSecretAwareCheckpointBeforeCleanup(t *testin runner := newBackupRunner(installation, true) deps := restoreTestDependencies(t, runner) fileErr := errors.New("cannot restore operator configuration") - deps.checkpoint = func(context.Context, config.Installation, CreateRequest) (Result, error) { + deps.checkpointLocked = func(context.Context, config.Installation, CreateRequest) (Result, error) { return Result{Path: "/tmp/recovery.zip"}, nil } var events []string @@ -791,7 +1095,7 @@ func TestRestoreReleasesBarrierOnlyAfterVerifiedRecoveryFromLostResponse(t *test } runner := &commandFailureRunner{fakeBackupRunner: backing, failures: []*commandFailure{&failure}} deps := restoreTestDependencies(t, runner) - deps.checkpoint = func(context.Context, config.Installation, CreateRequest) (Result, error) { + deps.checkpointLocked = func(context.Context, config.Installation, CreateRequest) (Result, error) { return Result{Path: "/tmp/recovery.zip"}, nil } deps.prepareRecovery = func(context.Context, config.Installation, string) (PreflightResult, error) { @@ -998,6 +1302,26 @@ func (runner *restoreFailureRunner) SessionInventoryScope() string { return runner.fakeBackupRunner.SessionInventoryScope() } +type lifecycleGateRunner struct { + *fakeBackupRunner + beforeFinalMaintenanceRelease func() +} + +func (runner *lifecycleGateRunner) Run(ctx context.Context, args []string, stdin io.Reader) (compose.Result, error) { + if strings.Contains(strings.Join(args, " "), "operator-command.js maintenance-deactivate") && runner.beforeFinalMaintenanceRelease != nil { + runner.beforeFinalMaintenanceRelease() + } + return runner.fakeBackupRunner.Run(ctx, args, stdin) +} + +func (runner *lifecycleGateRunner) Stream(ctx context.Context, args []string, stdin io.Reader, stdout io.Writer) (compose.Result, error) { + return runner.fakeBackupRunner.Stream(ctx, args, stdin, stdout) +} + +func (runner *lifecycleGateRunner) SessionInventoryScope() string { + return runner.fakeBackupRunner.SessionInventoryScope() +} + type fakeRestoreLock struct { release func() } @@ -1132,7 +1456,7 @@ func restoreTestDependencies(t *testing.T, runner archiveRunner) restoreDependen preflight: func(ctx context.Context, installation config.Installation, request PreflightRequest) (PreflightResult, error) { return Preflight(ctx, installation, request, permissivePreflightDependencies()) }, - checkpoint: func(context.Context, config.Installation, CreateRequest) (Result, error) { + checkpointLocked: func(context.Context, config.Installation, CreateRequest) (Result, error) { return Result{Path: "/tmp/default-checkpoint.zip"}, nil }, prepareRecovery: func(context.Context, config.Installation, string) (PreflightResult, error) {