package backup import ( "archive/tar" "bytes" "context" "errors" "fmt" "io" "os" "path/filepath" "strings" "testing" "time" "github.com/aritmolab/thothii/tools/tht/internal/authconfig" "github.com/aritmolab/thothii/tools/tht/internal/compose" "github.com/aritmolab/thothii/tools/tht/internal/config" "github.com/aritmolab/thothii/tools/tht/internal/lifecycle" "github.com/aritmolab/thothii/tools/tht/internal/safeio" ) func TestRestorePublicPathUsesConcreteProductionPreflight(t *testing.T) { root, err := filepath.EvalSymlinks(t.TempDir()) if err != nil { t.Fatal(err) } installation := config.Installation{ Path: filepath.Join(root, "deploy", "local-dev", "thothii-installation.yaml"), ProjectDirectory: root, } missing := filepath.Join(root, "missing.zip") _, err = Restore(context.Background(), installation, RestoreRequest{Archive: missing, Confirm: true}) if err == nil || strings.Contains(err.Error(), "dependencies are unavailable") || !strings.Contains(err.Error(), "backup archive") { t.Fatalf("Restore() error = %v, want production archive preflight", err) } } func TestVolumeRestoreCommandKeepsTarInputOpenWithoutTTY(t *testing.T) { args := volumeRestoreCommand("project_sessions") wantPrefix := []string{"run", "--rm", "--interactive", "--network", "none"} if len(args) < len(wantPrefix) || !equalStrings(args[:len(wantPrefix)], wantPrefix) { t.Fatalf("volumeRestoreCommand() prefix = %q, want %q", args, wantPrefix) } for _, arg := range args { if arg == "--tty" || arg == "-t" { t.Fatalf("volumeRestoreCommand() requests a TTY: %q", args) } } } func TestRestoreRestoresVerifiedVolumesInManifestOrderBeforeAuthenticationReset(t *testing.T) { installation := preflightTestInstallation(t) archive := filepath.Join(t.TempDir(), "restore-volumes.zip") sessionsTar := safeRestoreTar(t, "session.txt", "session") settingsTar := safeRestoreTar(t, "settings.json", "settings") writePreflightArchive(t, archive, preflightArchiveSpec{ volumes: []VolumeMetadata{ {LogicalName: "sessions", Name: "project_sessions", Driver: "local"}, {LogicalName: "settings", Name: "project_settings", Driver: "local"}, }, entries: []preflightArchiveEntry{ {path: "configuration/operator.env", body: []byte("safe")}, {path: "volumes/settings.tar", body: settingsTar, kind: EntryVolume, owner: "volume:settings", logicalName: "settings"}, {path: "volumes/sessions.tar", body: sessionsTar, kind: EntryVolume, owner: "volume:sessions", logicalName: "sessions"}, }, }) runner := newBackupRunner(installation, false) deps := restoreTestDependencies(t, runner) var events []string deps.restoreFile = func(_ context.Context, _ config.Installation, entry ArchiveEntryMetadata, _ io.Reader) error { events = append(events, "file:"+entry.Path) return nil } deps.restoreVolume = func(_ context.Context, _ config.Installation, volume VolumeMetadata, stream io.Reader) error { if _, err := io.ReadAll(stream); err != nil { return err } events = append(events, "volume:"+volume.LogicalName) return nil } deps.resetAuthenticationState = func(context.Context, config.Installation, archiveRunner) error { events = append(events, "reset-auth-state") return nil } if _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: archive, Confirm: true}, deps); err != nil { t.Fatal(err) } if got, want := events, []string{ "file:configuration/operator.env", "volume:sessions", "volume:settings", "reset-auth-state", }; !equalStrings(got, want) { t.Fatalf("restore events = %v, want %v", got, want) } } func safeRestoreTar(t *testing.T, name, contents string) []byte { t.Helper() var output bytes.Buffer writer := tar.NewWriter(&output) if err := writer.WriteHeader(&tar.Header{Name: name, Mode: 0o600, Size: int64(len(contents)), Typeflag: tar.TypeReg}); err != nil { t.Fatal(err) } if _, err := writer.Write([]byte(contents)); err != nil { t.Fatal(err) } if err := writer.Close(); err != nil { t.Fatal(err) } return output.Bytes() } func TestRestoreStoppedInstallationRunsCheckpointRestoreAndVerification(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) var events []string var checkpointRequest CreateRequest deps := restoreTestDependencies(t, runner) prepareRecovery := deps.prepareRecovery deps.checkpoint = func(_ context.Context, _ *lifecycle.Transaction, _ config.Installation, request CreateRequest) (Result, error) { events = append(events, "checkpoint") checkpointRequest = request return Result{Path: "/tmp/checkpoint.zip"}, nil } deps.prepareRecovery = func(ctx context.Context, target config.Installation, path string) (PreflightResult, error) { events = append(events, "prepare-recovery") return prepareRecovery(ctx, target, path) } deps.cleanupCheckpoint = func(string) error { events = append(events, "cleanup-checkpoint") return nil } deps.acquireTransaction = func(target config.Installation) (*lifecycle.Transaction, error) { events = append(events, "lock") return lifecycle.AcquireTransaction(target) } deps.restoreFile = func(_ context.Context, _ config.Installation, entry ArchiveEntryMetadata, _ io.Reader) error { events = append(events, "file:"+entry.Path) return nil } for _, name := range []string{"health", "doctor", "pi", "workspace"} { name := name deps.verify[name] = func(context.Context, config.Installation, archiveRunner) error { events = append(events, name) return nil } } result, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: archive, Confirm: true}, deps) if err != nil { t.Fatal(err) } if result.Checkpoint != "" || result.Restarted || !result.Verified { t.Fatalf("Restore() result = %#v", result) } if !checkpointRequest.IncludeSecrets || !checkpointRequest.Confirm { t.Fatalf("checkpoint request = %#v, want private confirmed secret-aware checkpoint", checkpointRequest) } if got, want := events, []string{"lock", "checkpoint", "prepare-recovery", "file:configuration/operator.env", "health", "doctor", "pi", "workspace", "cleanup-checkpoint"}; !equalStrings(got, want) { t.Fatalf("restore events = %v, want %v", got, want) } if err := lifecycleLockFreeAfterTerminalRestore(installation); err != nil { t.Fatal(err) } } func TestRestoreRejectsCombinedCandidateAndRecoveryStagingCapacityBeforeMutation(t *testing.T) { installation := preflightTestInstallation(t) candidateArchive := restoreArchive(t) recoveryArchive := restoreArchive(t) runner := newBackupRunner(installation, false) dependencies := restoreTestDependencies(t, runner) capacityChecks := 0 preflightDependencies := permissivePreflightDependencies() preflightDependencies.FreeBytes = func(string) (uint64, error) { capacityChecks++ if capacityChecks == 5 { return 0, nil } return 1 << 30, nil } dependencies.preflight = func(ctx context.Context, target config.Installation, request PreflightRequest) (PreflightResult, error) { return Preflight(ctx, target, request, preflightDependencies) } dependencies.checkpoint = func(context.Context, *lifecycle.Transaction, config.Installation, CreateRequest) (Result, error) { return Result{Path: recoveryArchive}, nil } dependencies.prepareRecovery = func(ctx context.Context, target config.Installation, path string) (PreflightResult, error) { return Preflight(ctx, target, PreflightRequest{Archive: path, Confirm: true, AllowExternalSecrets: true}, preflightDependencies) } mutated := false dependencies.restoreFile = func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error { mutated = true return nil } _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: candidateArchive, Confirm: true}, dependencies) if err == nil || !strings.Contains(err.Error(), "staging") { t.Fatalf("Restore() error = %v, want recovery staging capacity rejection", err) } if mutated { t.Fatal("restore mutated the installation before reserving candidate and recovery staging capacity") } if capacityChecks != 5 { t.Fatalf("free-space checks = %d, want candidate/recovery preflight and staging checks", capacityChecks) } } 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 } deps.acquireTransaction = lifecycle.AcquireTransaction if _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: archive, Confirm: true}, deps); err != nil { t.Fatal(err) } if target.archive != nil && target.archive.file != nil { if _, err := target.archive.file.Stat(); err == nil { t.Fatal("restore did not close the target archive") } } if err := lifecycleLockFreeAfterTerminalRestore(installation); err != nil { t.Fatal(err) } } func TestRestoreAcquiresLifecycleLockBeforeTargetDependentPreflight(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) deps.acquireTransaction = lifecycle.AcquireTransaction deps.preflight = func(ctx context.Context, target config.Installation, request PreflightRequest) (PreflightResult, error) { probe, err := lifecycle.Acquire(target) if err == nil { _ = probe.Release() return PreflightResult{}, errors.New("target-dependent preflight ran before lifecycle lock acquisition") } if !errors.Is(err, lifecycle.ErrLocked) { return PreflightResult{}, fmt.Errorf("probe lifecycle lock: %w", err) } return Preflight(ctx, target, request, permissivePreflightDependencies()) } if _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: archive, Confirm: true}, deps); err != nil { t.Fatalf("Restore() error = %v, want preflight protected by lifecycle lock", err) } if err := lifecycleLockFreeAfterTerminalRestore(installation); err != nil { t.Fatal(err) } } func TestRestoreStagesArchiveAfterCheckpointAndRejectsMutation(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("before")}}, }) runner := newBackupRunner(installation, false) deps := restoreTestDependencies(t, runner) deps.acquireTransaction = lifecycle.AcquireTransaction deps.checkpoint = func(context.Context, *lifecycle.Transaction, config.Installation, CreateRequest) (Result, error) { writePreflightArchive(t, archive, preflightArchiveSpec{ entries: []preflightArchiveEntry{{path: "configuration/operator.env", body: []byte("after!")}}, }) return Result{Path: filepath.Join(t.TempDir(), "checkpoint.zip")}, nil } var restored [][]byte deps.restoreFile = func(_ context.Context, _ config.Installation, _ ArchiveEntryMetadata, stream io.Reader) error { body, err := io.ReadAll(stream) if err != nil { return err } restored = append(restored, body) return nil } if _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: archive, Confirm: true}, deps); err == nil || !strings.Contains(err.Error(), "changed") { t.Fatalf("Restore() error = %v, want archive mutation refusal", err) } if len(restored) != 0 { t.Fatalf("restore applied mutated archive payloads: %q", restored) } if err := lifecycleLockFreeAfterTerminalRestore(installation); err != nil { t.Fatal(err) } } type restoreLifecycleTestOutcome struct { result RestoreResult err error } const restoreLifecycleTestWait = 2 * time.Second func stopRestoreLifecycleTimer(timer *time.Timer) { if !timer.Stop() { select { case <-timer.C: default: } } } func waitRestoreLifecycleSignal(ctx context.Context, signal <-chan struct{}, description string) error { timer := time.NewTimer(restoreLifecycleTestWait) defer stopRestoreLifecycleTimer(timer) select { case <-signal: return nil case <-ctx.Done(): return ctx.Err() case <-timer.C: return fmt.Errorf("timed out waiting for %s", description) } } func sendRestoreLifecycleSignal(ctx context.Context, signal chan<- struct{}, description string) error { timer := time.NewTimer(restoreLifecycleTestWait) defer stopRestoreLifecycleTimer(timer) select { case signal <- struct{}{}: return nil case <-ctx.Done(): return ctx.Err() case <-timer.C: return fmt.Errorf("timed out releasing %s", description) } } func waitRestoreLifecycleOutcome(ctx context.Context, done <-chan restoreLifecycleTestOutcome) (restoreLifecycleTestOutcome, bool, error) { timer := time.NewTimer(restoreLifecycleTestWait) defer stopRestoreLifecycleTimer(timer) select { case outcome := <-done: return outcome, true, nil case <-ctx.Done(): select { case outcome := <-done: return outcome, true, nil default: } return restoreLifecycleTestOutcome{}, false, ctx.Err() case <-timer.C: return restoreLifecycleTestOutcome{}, false, errors.New("timed out waiting for restore worker outcome") } } func waitRestoreLifecycleStage(ctx context.Context, stages <-chan string, done <-chan restoreLifecycleTestOutcome) (string, *restoreLifecycleTestOutcome, error) { timer := time.NewTimer(restoreLifecycleTestWait) defer stopRestoreLifecycleTimer(timer) select { case stage := <-stages: return stage, nil, nil case outcome := <-done: return "", &outcome, nil case <-ctx.Done(): select { case outcome := <-done: return "", &outcome, nil default: } return "", nil, ctx.Err() case <-timer.C: return "", nil, errors.New("timed out waiting for lifecycle stage") } } func waitRestoreError(ctx context.Context, done <-chan error) (error, bool, error) { timer := time.NewTimer(restoreLifecycleTestWait) defer stopRestoreLifecycleTimer(timer) select { case err := <-done: return err, true, nil case <-ctx.Done(): select { case err := <-done: return err, true, nil default: } return nil, false, ctx.Err() case <-timer.C: return nil, false, errors.New("timed out waiting for restore worker error") } } func waitRestoreAdmissionAttempt(ctx context.Context, attempt <-chan bool) (bool, error) { timer := time.NewTimer(restoreLifecycleTestWait) defer stopRestoreLifecycleTimer(timer) select { case allowed := <-attempt: return allowed, nil case <-ctx.Done(): return false, ctx.Err() case <-timer.C: return false, errors.New("timed out waiting for restore admission attempt") } } func releaseRestoreLifecycleGate(release chan<- struct{}) { select { case release <- struct{}{}: default: } } func joinRestoreLifecycleWorker(cancel context.CancelFunc, cancelGate context.CancelFunc, release chan<- struct{}, done <-chan restoreLifecycleTestOutcome) (restoreLifecycleTestOutcome, bool, error) { cancel() if cancelGate != nil { cancelGate() } releaseRestoreLifecycleGate(release) joinContext, stop := context.WithTimeout(context.Background(), restoreLifecycleTestWait) defer stop() return waitRestoreLifecycleOutcome(joinContext, done) } func joinRestoreErrorWorker(cancel context.CancelFunc, done <-chan error) (error, bool, error) { cancel() joinContext, stop := context.WithTimeout(context.Background(), restoreLifecycleTestWait) defer stop() return waitRestoreError(joinContext, done) } func releaseLifecycleStage(ctx context.Context, release chan<- struct{}, done <-chan restoreLifecycleTestOutcome) (restoreLifecycleTestOutcome, bool, error) { select { case outcome := <-done: return outcome, true, nil default: } timer := time.NewTimer(restoreLifecycleTestWait) defer stopRestoreLifecycleTimer(timer) select { case outcome := <-done: return outcome, true, nil case release <- struct{}{}: return restoreLifecycleTestOutcome{}, false, nil case <-ctx.Done(): select { case outcome := <-done: return outcome, true, nil default: } return restoreLifecycleTestOutcome{}, false, ctx.Err() case <-timer.C: return restoreLifecycleTestOutcome{}, false, errors.New("timed out releasing lifecycle stage") } } func TestReleaseLifecycleStageReturnsPrematureWorkerOutcome(t *testing.T) { want := errors.New("worker ended before the next lifecycle stage") done := make(chan restoreLifecycleTestOutcome, 1) done <- restoreLifecycleTestOutcome{err: want} got, terminal, err := releaseLifecycleStage(context.Background(), make(chan struct{}), done) if err != nil { t.Fatalf("releaseLifecycleStage() error = %v", err) } if !terminal || !errors.Is(got.err, want) { t.Fatalf("releaseLifecycleStage() = outcome %#v, terminal %t; want original worker error", got, terminal) } } func TestReleaseLifecycleStageHonorsCancellation(t *testing.T) { caller, cancel := context.WithCancel(context.Background()) cancel() started := time.Now() got, terminal, err := releaseLifecycleStage(caller, make(chan struct{}), make(chan restoreLifecycleTestOutcome)) if !errors.Is(err, context.Canceled) || terminal { t.Fatalf("releaseLifecycleStage() = outcome %#v, terminal %t, err %v; want prompt cancellation", got, terminal, err) } if elapsed := time.Since(started); elapsed >= time.Second { t.Fatalf("releaseLifecycleStage() cancellation took %s; want bounded prompt return", elapsed) } } func TestRestoreLifecycleCancellationJoinsWithWithheldGate(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{}) caller, cancel := context.WithCancel(context.Background()) defer cancel() gate := func(stage string) error { select { case stages <- stage: case <-caller.Done(): return caller.Err() } select { case <-continueStage: return nil case <-caller.Done(): return caller.Err() } } runner.beforeFinalMaintenanceRelease = func() { _ = gate("final-barrier-release") } deps := restoreTestDependencies(t, runner) deps.acquireTransaction = lifecycle.AcquireTransaction deps.checkpoint = func(context.Context, *lifecycle.Transaction, config.Installation, CreateRequest) (Result, error) { if err := gate("checkpoint"); err != nil { return Result{}, err } return Result{Path: filepath.Join(t.TempDir(), "checkpoint.zip")}, nil } done := make(chan restoreLifecycleTestOutcome, 1) go func() { result, err := restoreWithDependencies(caller, installation, RestoreRequest{Archive: archive, Confirm: true}, deps) done <- restoreLifecycleTestOutcome{result: result, err: err} }() var failures []error var outcome restoreLifecycleTestOutcome workerJoined := false stage, terminal, waitErr := waitRestoreLifecycleStage(caller, stages, done) if terminal != nil { outcome = *terminal workerJoined = true failures = append(failures, fmt.Errorf("restore ended before withheld checkpoint release: %w", outcome.err)) } else if waitErr != nil { failures = append(failures, fmt.Errorf("withheld checkpoint gate: %w", waitErr)) } else { if stage != "checkpoint" { failures = append(failures, fmt.Errorf("lifecycle stage = %q, want checkpoint", stage)) } diagnosticContext, stopDiagnostic := context.WithTimeout(caller, 100*time.Millisecond) _, diagnosticJoined, diagnosticErr := waitRestoreLifecycleOutcome(diagnosticContext, done) stopDiagnostic() if diagnosticJoined { failures = append(failures, errors.New("restore bypassed the intentionally withheld checkpoint release")) } else if !errors.Is(diagnosticErr, context.DeadlineExceeded) { failures = append(failures, fmt.Errorf("withheld checkpoint diagnostic = %w, want deadline", diagnosticErr)) } } if !workerJoined { outcome, workerJoined, waitErr = joinRestoreLifecycleWorker(cancel, cancel, continueStage, done) if waitErr != nil { failures = append(failures, fmt.Errorf("cancelled restore worker join: %w", waitErr)) } } if err := lifecycleLockFreeAfterTerminalRestore(installation); err != nil { failures = append(failures, err) } if !workerJoined { failures = append(failures, errors.New("restore worker outcome was never observed")) } else if !errors.Is(outcome.err, context.Canceled) { failures = append(failures, fmt.Errorf("restore cancellation outcome = %v, want context.Canceled", outcome.err)) } if len(failures) != 0 { t.Fatal(errors.Join(failures...)) } } 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) error, 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) error, _ context.CancelFunc) { deps.restoreFile = func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error { return gate("target") } deps.verify["health"] = func(context.Context, config.Installation, archiveRunner) error { return gate("verification") } }, }, { 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) error, _ context.CancelFunc) { deps.restoreFile = func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error { if err := gate("target-failure"); err != nil { return err } return targetFailure } deps.recover = func(context.Context, config.Installation, PreflightResult, *stagedArchive, bool, authProjectionRestoreTransaction) error { if err := gate("recovery"); err != nil { return err } 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) error, _ context.CancelFunc) { deps.restoreFile = func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error { if err := gate("target-failure"); err != nil { return err } return targetFailure } deps.recover = func(context.Context, config.Installation, PreflightResult, *stagedArchive, bool, authProjectionRestoreTransaction) error { if err := gate("recovery-failure"); err != nil { return err } 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) error, cancel context.CancelFunc) { deps.restoreFile = func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error { if err := gate("target-cancel"); err != nil { return err } cancel() return context.Canceled } deps.recover = func(context.Context, config.Installation, PreflightResult, *stagedArchive, bool, authProjectionRestoreTransaction) error { if err := gate("recovery"); err != nil { return err } 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)} recoveryArchive := filepath.Join(t.TempDir(), "recovery.zip") writePreflightArchive(t, recoveryArchive, preflightArchiveSpec{ installationID: fixture.installationID, entries: []preflightArchiveEntry{{path: "configuration/operator.env", body: []byte("checkpoint")}}, }) stages := make(chan string, 1) continueStage := make(chan struct{}, 1) caller, cancel := context.WithCancel(context.Background()) gateContext, cancelGate := context.WithCancel(context.Background()) defer cancel() defer cancelGate() gate := func(stage string) error { select { case stages <- stage: case <-gateContext.Done(): return gateContext.Err() } select { case <-continueStage: return nil case <-gateContext.Done(): return gateContext.Err() } } runner.beforeFinalMaintenanceRelease = func() { _ = gate("final-barrier-release") } deps := restoreTestDependencies(t, runner) deps.prepareRecovery = func(ctx context.Context, target config.Installation, _ string) (PreflightResult, error) { return Preflight(ctx, target, PreflightRequest{Archive: recoveryArchive, Confirm: true, AllowExternalSecrets: true}, permissivePreflightDependencies()) } deps.acquireTransaction = lifecycle.AcquireTransaction deps.checkpoint = func(context.Context, *lifecycle.Transaction, config.Installation, CreateRequest) (Result, error) { if err := gate("checkpoint"); err != nil { return Result{}, err } return Result{Path: filepath.Join(t.TempDir(), "checkpoint.zip")}, nil } deps.cleanupCheckpoint = func(string) error { return gate("checkpoint-cleanup") } scenario.configure(&deps, runner, gate, cancel) done := make(chan restoreLifecycleTestOutcome, 1) go func() { result, err := restoreWithDependencies(caller, installation, RestoreRequest{Archive: archive, Confirm: true}, deps) done <- restoreLifecycleTestOutcome{result: result, err: err} }() var failures []error var terminal *restoreLifecycleTestOutcome workerJoined := false aborted := false for stageIndex, wantStage := range scenario.stages { stage, outcome, waitErr := waitRestoreLifecycleStage(gateContext, stages, done) if outcome != nil { terminal = outcome workerJoined = true failures = append(failures, fmt.Errorf("restore ended before lifecycle stage %q: %w", wantStage, outcome.err)) break } if waitErr != nil { failures = append(failures, fmt.Errorf("lifecycle stage %q: %w", wantStage, waitErr)) aborted = true break } if stage != wantStage { failures = append(failures, fmt.Errorf("lifecycle stage = %q, want %q", stage, wantStage)) } if err := competingRestoreAndBackupEntry(installation, archive, t); err != nil { failures = append(failures, fmt.Errorf("%s: %w", wantStage, err)) } releasedOutcome, workerDone, err := releaseLifecycleStage(gateContext, continueStage, done) if err != nil { failures = append(failures, fmt.Errorf("release lifecycle stage %q: %w", wantStage, err)) aborted = true break } if workerDone { terminal = &releasedOutcome workerJoined = true if stageIndex+1 < len(scenario.stages) { failures = append(failures, fmt.Errorf("restore ended before lifecycle stage %q: %w", scenario.stages[stageIndex+1], releasedOutcome.err)) } break } } var result restoreLifecycleTestOutcome if terminal != nil { result = *terminal } else if aborted { var joinErr error result, workerJoined, joinErr = joinRestoreLifecycleWorker(cancel, cancelGate, continueStage, done) if joinErr != nil { failures = append(failures, joinErr) } } else { var waitErr error result, workerJoined, waitErr = waitRestoreLifecycleOutcome(context.Background(), done) if waitErr != nil { failures = append(failures, waitErr) result, workerJoined, waitErr = joinRestoreLifecycleWorker(cancel, cancelGate, continueStage, done) if waitErr != nil { failures = append(failures, waitErr) } } } 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 !workerJoined { failures = append(failures, errors.New("restore worker outcome was never joined")) } 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.acquireTransaction = lifecycle.AcquireTransaction restoreDeps.checkpoint = func(context.Context, *lifecycle.Transaction, 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{}) firstContext, firstCancel := context.WithCancel(context.Background()) defer firstCancel() firstRunner := newBackupRunner(installation, true) firstDeps := restoreTestDependencies(t, firstRunner) recoveryArchive := filepath.Join(t.TempDir(), "first-recovery.zip") writePreflightArchive(t, recoveryArchive, preflightArchiveSpec{ installationID: fixture.installationID, entries: []preflightArchiveEntry{{path: "configuration/operator.env", body: []byte("checkpoint")}}, }) firstDeps.prepareRecovery = func(ctx context.Context, target config.Installation, _ string) (PreflightResult, error) { return Preflight(ctx, target, PreflightRequest{Archive: recoveryArchive, Confirm: true, AllowExternalSecrets: true}, permissivePreflightDependencies()) } firstDeps.acquireTransaction = lifecycle.AcquireTransaction firstDeps.checkpoint = func(ctx context.Context, _ *lifecycle.Transaction, _ config.Installation, _ CreateRequest) (Result, error) { checkpointState = targetState close(checkpointEntered) select { case <-continueCheckpoint: case <-ctx.Done(): return Result{}, ctx.Err() } 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, *stagedArchive, bool, authProjectionRestoreTransaction) error { recoveryObserved = targetState targetState = checkpointState firstRunner.running, firstRunner.coreRunning = true, true return nil } done := make(chan error, 1) go func() { _, err := restoreWithDependencies(firstContext, installation, RestoreRequest{Archive: archive, Confirm: true}, firstDeps) done <- err }() if err := waitRestoreLifecycleSignal(context.Background(), checkpointEntered, "first restore recovery checkpoint"); err != nil { _, _, joinErr := joinRestoreErrorWorker(firstCancel, done) if joinErr != nil { err = errors.Join(err, joinErr) } if lockErr := lifecycleLockFreeAfterTerminalRestore(installation); lockErr != nil { err = errors.Join(err, lockErr) } t.Fatal(err) } interleavedCheckpoint := false interleavedMutation := false interleavedDeps := restoreTestDependencies(t, newBackupRunner(installation, false)) interleavedDeps.acquireTransaction = lifecycle.AcquireTransaction interleavedDeps.checkpoint = func(context.Context, *lifecycle.Transaction, 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) if err := sendRestoreLifecycleSignal(context.Background(), continueCheckpoint, "first restore checkpoint"); err != nil { _, _, joinErr := joinRestoreErrorWorker(firstCancel, done) if joinErr != nil { err = errors.Join(err, joinErr) } if lockErr := lifecycleLockFreeAfterTerminalRestore(installation); lockErr != nil { err = errors.Join(err, lockErr) } t.Fatal(err) } firstErr, joined, waitErr := waitRestoreError(context.Background(), done) if waitErr != nil { firstErr, joined, waitErr = joinRestoreErrorWorker(firstCancel, done) } if waitErr != nil { if lockErr := lifecycleLockFreeAfterTerminalRestore(installation); lockErr != nil { waitErr = errors.Join(waitErr, lockErr) } t.Fatal(waitErr) } if !joined { t.Fatal("first restore worker outcome was never joined") } 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) runner := newBackupRunner(installation, true) deps := restoreTestDependencies(t, runner) var events []string deps.restoreFile = func(_ context.Context, _ config.Installation, entry ArchiveEntryMetadata, _ io.Reader) error { events = append(events, "file:"+entry.Path) return nil } deps.resetAuthenticationState = func(context.Context, config.Installation, archiveRunner) error { events = append(events, "reset-auth-state") return nil } for _, name := range []string{"health", "doctor", "pi", "workspace"} { name := name deps.verify[name] = func(context.Context, config.Installation, archiveRunner) error { events = append(events, name) return nil } } result, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: archive, Confirm: true}, deps) if err != nil { t.Fatal(err) } if !result.Restarted || !result.Verified { t.Fatalf("restore result = %#v", result) } if got, want := events, []string{"file:configuration/operator.env", "reset-auth-state", "health", "doctor", "pi", "workspace"}; !equalStrings(got, want) { t.Fatalf("restore events = %v, want %v", got, want) } } func TestResetAuthenticationStateCreatesOnlyPrivateEmptyStateDirectories(t *testing.T) { installation := preflightTestInstallation(t) runner := &authenticationStateResetRunner{} if err := resetAuthenticationState(context.Background(), installation, runner); err != nil { t.Fatal(err) } joined := strings.Join(runner.args, "\x00") for _, required := range []string{ "run", "--rm", "--no-deps", "--no-TTY", "--user", "0:0", "--entrypoint", "sh", "core", "-ceu", "test ! -L /data/auth", "install -d -o 10001 -g 10001 -m 0700 /data/auth", "find /data/auth -mindepth 1 -maxdepth 1 -exec rm -rf -- {} +", "install -d -o 10001 -g 10001 -m 0700 /data/auth/sessions /data/auth/oidc", "find /data/auth/sessions /data/auth/oidc -mindepth 1 -print -quit", } { if !strings.Contains(joined, required) { t.Fatalf("authentication state reset command omits %q: %#v", required, runner.args) } } } func TestResetAuthenticationStateFailsClosedForDamagedRoot(t *testing.T) { installation := preflightTestInstallation(t) runner := &authenticationStateResetRunner{result: compose.Result{ExitCode: 1}, err: errors.New("damaged auth root")} err := resetAuthenticationState(context.Background(), installation, runner) if err == nil || !strings.Contains(err.Error(), "reset authentication state") { t.Fatalf("resetAuthenticationState() error = %v, want bounded fail-closed error", err) } } func TestRestorePreflightFailureDoesNotMutateTarget(t *testing.T) { installation := preflightTestInstallation(t) runner := newBackupRunner(installation, true) deps := restoreTestDependencies(t, runner) preflightErr := errors.New("archive cannot be restored") deps.preflight = func(context.Context, config.Installation, PreflightRequest) (PreflightResult, error) { return PreflightResult{}, preflightErr } checkpointCalls := 0 deps.checkpoint = func(context.Context, *lifecycle.Transaction, config.Installation, CreateRequest) (Result, error) { checkpointCalls++ return Result{}, nil } restoredFiles := 0 deps.restoreFile = func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error { restoredFiles++ return nil } result, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: "unreadable.zip", Confirm: true}, deps) if !errors.Is(err, preflightErr) { t.Fatalf("restore error = %v, want preflight error", err) } if result != (RestoreResult{}) || checkpointCalls != 0 || restoredFiles != 0 || runner.stopCount != 0 || runner.startCount != 0 || !runner.running { t.Fatalf("preflight failure mutated target: result=%#v checkpoint=%d files=%d stops=%d starts=%d running=%t", result, checkpointCalls, restoredFiles, runner.stopCount, runner.startCount, runner.running) } } func TestRestoreCheckpointFailureDoesNotMutateTarget(t *testing.T) { installation := preflightTestInstallation(t) archive := restoreArchive(t) runner := newBackupRunner(installation, true) deps := restoreTestDependencies(t, runner) checkpointErr := errors.New("checkpoint unavailable") deps.checkpoint = func(context.Context, *lifecycle.Transaction, config.Installation, CreateRequest) (Result, error) { return Result{}, checkpointErr } restoredFiles := 0 deps.restoreFile = func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error { restoredFiles++ return nil } result, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: archive, Confirm: true}, deps) if !errors.Is(err, checkpointErr) { t.Fatalf("restore error = %v, want checkpoint error", err) } if result != (RestoreResult{}) || restoredFiles != 0 || runner.stopCount != 0 || runner.startCount != 0 || !runner.running { t.Fatalf("checkpoint failure mutated target: result=%#v files=%d stops=%d starts=%d running=%t", result, restoredFiles, runner.stopCount, runner.startCount, runner.running) } } func TestRestoreFileFailureRollsBackSecretAwareCheckpointBeforeCleanup(t *testing.T) { installation := preflightTestInstallation(t) archive := restoreArchive(t) runner := newBackupRunner(installation, true) deps := restoreTestDependencies(t, runner) fileErr := errors.New("cannot restore operator configuration") deps.checkpoint = func(context.Context, *lifecycle.Transaction, config.Installation, CreateRequest) (Result, error) { return Result{Path: "/tmp/recovery.zip"}, nil } var events []string deps.recover = func(context.Context, config.Installation, PreflightResult, *stagedArchive, bool, authProjectionRestoreTransaction) error { events = append(events, "recover") return nil } deps.cleanupCheckpoint = func(string) error { events = append(events, "cleanup") return nil } deps.restoreFile = func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error { events = append(events, "mutate") return fileErr } result, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: archive, Confirm: true}, deps) if !errors.Is(err, fileErr) { t.Fatalf("restore error = %v, want file error", err) } if result.Checkpoint != "" { t.Fatalf("recovery checkpoint = %q, want no retained secret-bearing path", result.Checkpoint) } if got, want := events, []string{"mutate", "recover", "cleanup"}; !equalStrings(got, want) { t.Fatalf("failure recovery events = %v, want %v", got, want) } } func TestRestoreFailureAfterAuthenticationMutationRollsBackAndClearsRuntimeState(t *testing.T) { installation := preflightTestInstallation(t) archive := restoreArchive(t) runner := newBackupRunner(installation, false) deps := restoreTestDependencies(t, runner) resetErr := errors.New("authentication reset failed") var events []string deps.restoreFile = func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error { events = append(events, "auth-config-mutated") return nil } deps.resetAuthenticationState = func(context.Context, config.Installation, archiveRunner) error { events = append(events, "auth-runtime-reset-failed") return resetErr } deps.recover = func(context.Context, config.Installation, PreflightResult, *stagedArchive, bool, authProjectionRestoreTransaction) error { events = append(events, "secret-aware-recovery-and-reauth-reset") return nil } deps.cleanupCheckpoint = func(string) error { events = append(events, "checkpoint-cleaned") return nil } _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: archive, Confirm: true}, deps) if !errors.Is(err, resetErr) { t.Fatalf("restore error = %v, want reset failure", err) } if got, want := events, []string{"auth-config-mutated", "auth-runtime-reset-failed", "secret-aware-recovery-and-reauth-reset", "checkpoint-cleaned"}; !equalStrings(got, want) { t.Fatalf("post-auth recovery events = %v, want %v", got, want) } } func TestRestoreCleanupFailureDoesNotSuppressRollback(t *testing.T) { installation := preflightTestInstallation(t) archive := restoreArchive(t) deps := restoreTestDependencies(t, newBackupRunner(installation, false)) mutationErr := errors.New("post-mutation failure") cleanupErr := errors.New("checkpoint cleanup failure") recovered := false deps.restoreFile = func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error { return mutationErr } deps.recover = func(context.Context, config.Installation, PreflightResult, *stagedArchive, bool, authProjectionRestoreTransaction) error { recovered = true return nil } deps.cleanupCheckpoint = func(string) error { return cleanupErr } _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: archive, Confirm: true}, deps) if !recovered || !errors.Is(err, mutationErr) || !errors.Is(err, cleanupErr) { t.Fatalf("restore error = %v recovered=%t, want joined mutation/cleanup failure after rollback", err, recovered) } } func TestRecoveryCheckpointCleanupRejectsSymlinksAndPermissiveFiles(t *testing.T) { root := t.TempDir() target := filepath.Join(root, "checkpoint.zip") if err := os.WriteFile(target, []byte("secret"), 0o600); err != nil { t.Fatal(err) } link := filepath.Join(root, "checkpoint-link.zip") if err := os.Symlink(target, link); err != nil { t.Skipf("symlink unavailable: %v", err) } if err := cleanupRecoveryCheckpoint(link); err == nil { t.Fatal("cleanupRecoveryCheckpoint accepted a symlink") } if _, err := os.Stat(target); err != nil { t.Fatalf("symlink target was changed: %v", err) } if err := os.Chmod(target, 0o644); err != nil { t.Fatal(err) } if err := cleanupRecoveryCheckpoint(target); err == nil { t.Fatal("cleanupRecoveryCheckpoint accepted a permissive secret checkpoint") } } func TestRecoveryCheckpointCleanupRemovesOnlyPrivateRegularFile(t *testing.T) { root, err := filepath.EvalSymlinks(t.TempDir()) if err != nil { t.Fatal(err) } checkpoint := filepath.Join(root, "checkpoint.zip") if err := os.WriteFile(checkpoint, []byte("secret checkpoint"), 0o600); err != nil { t.Fatal(err) } if err := safeio.ProtectPrivateRegular(checkpoint); err != nil { t.Fatal(err) } if err := cleanupRecoveryCheckpoint(checkpoint); err != nil { t.Fatal(err) } if _, err := os.Lstat(checkpoint); !errors.Is(err, os.ErrNotExist) { t.Fatalf("private checkpoint still exists after cleanup: %v", err) } } func TestRestoreStartFailureRecoversPreviouslyRunningTarget(t *testing.T) { installation := preflightTestInstallation(t) archive := restoreArchive(t) backingRunner := newBackupRunner(installation, true) deps := restoreTestDependencies(t, failStartRestoreRunner{fakeBackupRunner: backingRunner}) recovered := false deps.recover = func(context.Context, config.Installation, PreflightResult, *stagedArchive, bool, authProjectionRestoreTransaction) error { recovered = true backingRunner.running = true return nil } _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: archive, Confirm: true}, deps) if err == nil || !strings.Contains(err.Error(), "start refused") { t.Fatalf("restore error = %v, want restart failure", err) } if !recovered || !backingRunner.running { t.Fatalf("failed restart was not rolled back: recovered=%t running=%t", recovered, backingRunner.running) } } func TestRestoreVerificationFailureRecoversPreviouslyRunningTarget(t *testing.T) { installation := preflightTestInstallation(t) archive := restoreArchive(t) runner := newBackupRunner(installation, true) deps := restoreTestDependencies(t, runner) verificationErr := errors.New("Pi is unavailable") deps.verify["pi"] = func(context.Context, config.Installation, archiveRunner) error { return verificationErr } recovered := false deps.recover = func(context.Context, config.Installation, PreflightResult, *stagedArchive, bool, authProjectionRestoreTransaction) error { recovered = true return nil } _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: archive, Confirm: true}, deps) if !errors.Is(err, verificationErr) { t.Fatalf("restore error = %v, want verification failure", err) } if !recovered || !runner.running { t.Fatalf("verification failure was not rolled back: recovered=%t running=%t", recovered, runner.running) } } func TestRestoreKeepsAdmissionBarrierActiveUntilVerificationCommits(t *testing.T) { installation := preflightTestInstallation(t) archive := restoreArchive(t) runner := newBackupRunner(installation, true) deps := restoreTestDependencies(t, runner) verificationEntered := make(chan struct{}) allowVerification := make(chan struct{}) caller, cancel := context.WithCancel(context.Background()) defer cancel() deps.verify["health"] = func(ctx context.Context, _ config.Installation, _ archiveRunner) error { close(verificationEntered) select { case <-allowVerification: case <-ctx.Done(): return ctx.Err() } if !runner.maintenance { return errors.New("maintenance barrier was removed before verification committed") } return nil } done := make(chan restoreLifecycleTestOutcome, 1) go func() { result, err := restoreWithDependencies(caller, installation, RestoreRequest{Archive: archive, Confirm: true}, deps) done <- restoreLifecycleTestOutcome{result: result, err: err} }() var failures []error if err := waitRestoreLifecycleSignal(context.Background(), verificationEntered, "post-start verification"); err != nil { failures = append(failures, err) } // A newly admitted operation observes the same durable barrier as the backend gate. No // operation may enter after restore mutation and before the verification transaction commits. admissionAttempt := make(chan bool, 1) go func() { admissionAttempt <- !runner.maintenance }() admissionAllowed, err := waitRestoreAdmissionAttempt(context.Background(), admissionAttempt) if err != nil { failures = append(failures, err) } if err := sendRestoreLifecycleSignal(context.Background(), allowVerification, "post-start verification"); err != nil { failures = append(failures, err) } outcome, workerJoined, err := waitRestoreLifecycleOutcome(context.Background(), done) if err != nil { failures = append(failures, err) } if !workerJoined { joinedOutcome, joined, joinErr := joinRestoreLifecycleWorker(cancel, nil, nil, done) if joined { outcome = joinedOutcome workerJoined = true } if joinErr != nil { failures = append(failures, joinErr) } } if admissionAllowed { failures = append(failures, errors.New("a new operation could enter while restore verification was still in progress")) } if outcome.err != nil || !outcome.result.Verified { failures = append(failures, fmt.Errorf("restore result = %#v, %v; want successful verified restore", outcome.result, outcome.err)) } if runner.maintenance { failures = append(failures, errors.New("maintenance barrier remained active after successful verification")) } if err := lifecycleLockFreeAfterTerminalRestore(installation); err != nil { failures = append(failures, err) } if !workerJoined { failures = append(failures, errors.New("restore worker outcome was never joined")) } if len(failures) != 0 { t.Fatal(errors.Join(failures...)) } } func TestRestoreRecoversBehindBarrierForEveryVerificationFailure(t *testing.T) { for _, verification := range []string{"health", "doctor", "pi", "workspace"} { verification := verification t.Run(verification, func(t *testing.T) { installation := preflightTestInstallation(t) runner := newBackupRunner(installation, true) deps := restoreTestDependencies(t, runner) verificationErr := fmt.Errorf("%s verification failed", verification) deps.verify[verification] = func(context.Context, config.Installation, archiveRunner) error { return verificationErr } var recoveryBarrierActive bool var recoveryContext cleanupContextObservation deps.recover = func(ctx context.Context, _ config.Installation, _ PreflightResult, _ *stagedArchive, _ bool, _ authProjectionRestoreTransaction) error { recoveryContext = observeCleanupContext(ctx) recoveryBarrierActive = runner.maintenance runner.running, runner.coreRunning = true, true return nil } _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: restoreArchive(t), Confirm: true}, deps) if !errors.Is(err, verificationErr) { t.Fatalf("Restore() error = %v, want %v", err, verificationErr) } if !recoveryBarrierActive { t.Fatal("recovery started after the admission barrier had been removed") } assertIndependentBoundedCleanupContexts(t, []cleanupContextObservation{recoveryContext}) if runner.maintenance || !runner.running { t.Fatalf("recovery cleanup state = maintenance:%t running:%t", runner.maintenance, runner.running) } }) } } func TestRestoreDoesNotRollbackAfterFinalDeactivationResponseLoss(t *testing.T) { installation := preflightTestInstallation(t) backing := newBackupRunner(installation, true) deactivationResponseLost := errors.New("deactivation response lost") runner := &commandFailureRunner{ fakeBackupRunner: backing, failures: []*commandFailure{{ match: func(command string) bool { return strings.Contains(command, " maintenance-deactivate") }, err: deactivationResponseLost, effect: func() { // The barrier was removed, but Docker lost the response to the operator command. backing.maintenance = false }, remaining: 1, }}, } deps := restoreTestDependencies(t, runner) recoveryCalls := 0 deps.recover = func(context.Context, config.Installation, PreflightResult, *stagedArchive, bool, authProjectionRestoreTransaction) error { recoveryCalls++ return nil } _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: restoreArchive(t), Confirm: true}, deps) if !errors.Is(err, deactivationResponseLost) { t.Fatalf("Restore() error = %v, want %v", err, deactivationResponseLost) } if recoveryCalls != 0 { t.Fatalf("successful verification triggered an unnecessary rollback: %d recoveries", recoveryCalls) } if backing.maintenance || !backing.running { t.Fatalf("final deactivation cleanup state = maintenance:%t running:%t", backing.maintenance, backing.running) } } func TestRestoreUsesBoundedIndependentCleanupContextsAfterCanceledStopResponse(t *testing.T) { installation := preflightTestInstallation(t) backing := newBackupRunner(installation, true) caller, cancel := context.WithCancel(context.Background()) t.Cleanup(cancel) runner := &commandFailureRunner{ fakeBackupRunner: backing, failures: []*commandFailure{{ match: func(command string) bool { return strings.HasSuffix(command, " stop") }, err: context.Canceled, effect: func() { backing.stopCount++ backing.running, backing.coreRunning = false, false cancel() }, remaining: 1, }}, } _, err := restoreWithDependencies(caller, installation, RestoreRequest{Archive: restoreArchive(t), Confirm: true}, restoreTestDependencies(t, runner)) if !errors.Is(err, context.Canceled) { t.Fatalf("Restore() error = %v, want context cancellation", err) } if backing.maintenance || !backing.running || backing.startCount != 1 { t.Fatalf("cancelled cleanup state = maintenance:%t running:%t starts:%d", backing.maintenance, backing.running, backing.startCount) } assertIndependentBoundedCleanupContexts(t, runner.cleanupCommandContexts) } func TestRestoreUsesBoundedRecoveryContextAfterPostMutationCancellation(t *testing.T) { installation := preflightTestInstallation(t) runner := newBackupRunner(installation, true) caller, cancel := context.WithCancel(context.Background()) t.Cleanup(cancel) deps := restoreTestDependencies(t, runner) deps.restoreFile = func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error { cancel() return context.Canceled } var recoveryContext cleanupContextObservation var recoveryBarrierActive bool deps.recover = func(ctx context.Context, _ config.Installation, _ PreflightResult, _ *stagedArchive, _ bool, _ authProjectionRestoreTransaction) error { recoveryContext = observeCleanupContext(ctx) recoveryBarrierActive = runner.maintenance runner.running, runner.coreRunning = true, true return nil } _, err := restoreWithDependencies(caller, installation, RestoreRequest{Archive: restoreArchive(t), Confirm: true}, deps) if !errors.Is(err, context.Canceled) { t.Fatalf("Restore() error = %v, want cancelled post-mutation restore", err) } if !recoveryBarrierActive { t.Fatal("post-mutation recovery started after the admission barrier was removed") } assertIndependentBoundedCleanupContexts(t, []cleanupContextObservation{recoveryContext}) if runner.maintenance || !runner.running { t.Fatalf("post-mutation cancellation cleanup state = maintenance:%t running:%t", runner.maintenance, runner.running) } } func TestRestoreJoinsPrimaryRestartAndDeactivationFailures(t *testing.T) { installation := preflightTestInstallation(t) backing := newBackupRunner(installation, true) stopResponseLost := errors.New("stop response lost") startResponseLost := errors.New("restart response lost") deactivationResponseLost := errors.New("deactivation response lost") runner := &commandFailureRunner{ fakeBackupRunner: backing, failures: []*commandFailure{ { match: func(command string) bool { return strings.HasSuffix(command, " stop") }, err: stopResponseLost, effect: func() { backing.stopCount++ backing.running, backing.coreRunning = false, false }, remaining: 1, }, { match: func(command string) bool { return strings.HasSuffix(command, " start") }, err: startResponseLost, effect: func() { backing.startCount++ backing.running, backing.coreRunning = true, true }, remaining: 1, }, { match: func(command string) bool { return strings.Contains(command, " maintenance-deactivate") }, err: deactivationResponseLost, effect: func() { backing.maintenance = false }, remaining: 1, }, }, } _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: restoreArchive(t), Confirm: true}, restoreTestDependencies(t, runner)) if !errors.Is(err, stopResponseLost) || !errors.Is(err, startResponseLost) || !errors.Is(err, deactivationResponseLost) { t.Fatalf("Restore() error = %v, want joined stop/restart/deactivation failures", err) } if backing.maintenance || !backing.running { t.Fatalf("response-loss cleanup state = maintenance:%t running:%t", backing.maintenance, backing.running) } } func TestRecoverRestoreTransactionVerifiesRecoveredStateBeforeReturning(t *testing.T) { installation := preflightTestInstallation(t) recovery, err := Preflight(context.Background(), installation, PreflightRequest{Archive: restoreArchive(t), Confirm: true, AllowExternalSecrets: true}, permissivePreflightDependencies()) if err != nil { t.Fatal(err) } defer recovery.CloseArchive() runner := newBackupRunner(installation, true) deps := restoreTestDependencies(t, runner) var checks []string for _, name := range []string{"health", "doctor", "pi", "workspace"} { name := name deps.verify[name] = func(context.Context, config.Installation, archiveRunner) error { checks = append(checks, name) return nil } } staged := stageRecoveryForTest(t, installation, recovery) if err := recoverRestoreTransaction(context.Background(), installation, recovery, staged, true, deps, nil); err != nil { t.Fatal(err) } if got, want := checks, []string{"health", "doctor", "pi", "workspace"}; !equalStrings(got, want) { t.Fatalf("recovery checks = %v, want %v", got, want) } } func TestRecoverRestoreTransactionFailsClosedForEveryVerification(t *testing.T) { for _, verification := range []string{"health", "doctor", "pi", "workspace"} { verification := verification t.Run(verification, func(t *testing.T) { installation := preflightTestInstallation(t) recovery, err := Preflight(context.Background(), installation, PreflightRequest{Archive: restoreArchive(t), Confirm: true, AllowExternalSecrets: true}, permissivePreflightDependencies()) if err != nil { t.Fatal(err) } defer recovery.CloseArchive() runner := newBackupRunner(installation, true) runner.maintenance = true deps := restoreTestDependencies(t, runner) verificationErr := fmt.Errorf("recovery %s verification failed", verification) deps.verify[verification] = func(context.Context, config.Installation, archiveRunner) error { if !runner.maintenance { t.Fatal("recovery verification ran after the maintenance barrier was removed") } return verificationErr } staged := stageRecoveryForTest(t, installation, recovery) err = recoverRestoreTransaction(context.Background(), installation, recovery, staged, true, deps, nil) if !errors.Is(err, verificationErr) { t.Fatalf("recoverRestoreTransaction() error = %v, want %v", err, verificationErr) } if !runner.maintenance { t.Fatal("recovery removed the maintenance barrier after a failed verification") } }) } } func TestRestoreReleasesBarrierOnlyAfterVerifiedRecoveryFromLostResponse(t *testing.T) { for _, scenario := range []struct { name string failure commandFailure }{ { name: "recovery stop", failure: commandFailure{ match: func(command string) bool { return strings.HasSuffix(command, " stop") }, err: errors.New("recovery stop response lost"), effect: func() { // The stop took effect before Docker lost the response. }, skip: 1, remaining: 1, }, }, { name: "recovery start", failure: commandFailure{ match: func(command string) bool { return strings.HasSuffix(command, " start") }, err: errors.New("recovery start response lost"), effect: func() { // The start took effect before Docker lost the response. }, remaining: 1, }, }, } { scenario := scenario t.Run(scenario.name, func(t *testing.T) { installation := preflightTestInstallation(t) archive := restoreArchive(t) recovery, err := Preflight(context.Background(), installation, PreflightRequest{Archive: archive, Confirm: true, AllowExternalSecrets: true}, permissivePreflightDependencies()) if err != nil { t.Fatal(err) } backing := newBackupRunner(installation, true) failure := scenario.failure originalEffect := failure.effect failure.effect = func() { if strings.Contains(scenario.name, "stop") { backing.stopCount++ backing.running, backing.coreRunning = false, false } else { backing.startCount++ backing.running, backing.coreRunning = true, true } if originalEffect != nil { originalEffect() } } runner := &commandFailureRunner{fakeBackupRunner: backing, failures: []*commandFailure{&failure}} deps := restoreTestDependencies(t, runner) deps.checkpoint = func(context.Context, *lifecycle.Transaction, config.Installation, CreateRequest) (Result, error) { return Result{Path: "/tmp/recovery.zip"}, nil } deps.prepareRecovery = func(context.Context, config.Installation, string) (PreflightResult, error) { return recovery, nil } restoreCalls := 0 mutationErr := errors.New("target mutation failed") deps.restoreFile = func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error { restoreCalls++ if restoreCalls == 1 { return mutationErr } return nil } deps.recover = func(ctx context.Context, target config.Installation, checkpoint PreflightResult, staged *stagedArchive, wasRunning bool, transaction authProjectionRestoreTransaction) error { return recoverRestoreTransaction(ctx, target, checkpoint, staged, wasRunning, deps, transaction) } _, err = restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: archive, Confirm: true}, deps) if !errors.Is(err, mutationErr) || !errors.Is(err, failure.err) { t.Fatalf("Restore() error = %v, want joined mutation and recovery response-loss errors", err) } if backing.maintenance || !backing.running { t.Fatalf("verified recovery state = maintenance:%t running:%t", backing.maintenance, backing.running) } assertIndependentBoundedCleanupContexts(t, runner.cleanupCommandContexts) }) } } func TestRestoreRefusesActiveSessionsWithoutDrain(t *testing.T) { installation := preflightTestInstallation(t) archive := restoreArchive(t) runner := newBackupRunner(installation, true) runner.sessionResponses = []string{`[{"status":"running","archived":false}]`} deps := restoreTestDependencies(t, runner) restoredFiles := 0 deps.restoreFile = func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error { restoredFiles++ return nil } _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: archive, Confirm: true}, deps) if !errors.Is(err, ErrActiveSessions) { t.Fatalf("restore error = %v, want active-session refusal", err) } if restoredFiles != 0 || runner.stopCount != 0 || runner.startCount != 0 || !runner.running || runner.maintenance { t.Fatalf("active-session refusal changed lifecycle state: files=%d stops=%d starts=%d running=%t maintenance=%t", restoredFiles, runner.stopCount, runner.startCount, runner.running, runner.maintenance) } } func TestRestoreCleansMaintenanceAfterActivationFailure(t *testing.T) { installation := preflightTestInstallation(t) backing := newBackupRunner(installation, true) activationErr := errors.New("activation response lost") runner := &restoreFailureRunner{ fakeBackupRunner: backing, failContains: "operator-command.js maintenance-activate", err: activationErr, beforeFailure: func() { backing.maintenance = true }, } deps := restoreTestDependencies(t, runner) _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: restoreArchive(t), Confirm: true}, deps) if !errors.Is(err, activationErr) { t.Fatalf("restore error = %v, want activation failure", err) } if backing.maintenance || !backing.running || runner.matchCount != 1 { t.Fatalf("activation cleanup state: maintenance=%t running=%t matches=%d", backing.maintenance, backing.running, runner.matchCount) } } func TestRestoreRestartsAndCleansMaintenanceAfterStopFailure(t *testing.T) { installation := preflightTestInstallation(t) backing := newBackupRunner(installation, true) stopErr := errors.New("stop response lost") runner := &restoreFailureRunner{ fakeBackupRunner: backing, failSuffix: " stop", err: stopErr, beforeFailure: func() { backing.stopCount++ backing.running, backing.coreRunning = false, false }, } deps := restoreTestDependencies(t, runner) _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: restoreArchive(t), Confirm: true}, deps) if !errors.Is(err, stopErr) { t.Fatalf("restore error = %v, want stop failure", err) } if backing.maintenance || !backing.running || backing.startCount != 1 { t.Fatalf("stop cleanup state: maintenance=%t running=%t starts=%d", backing.maintenance, backing.running, backing.startCount) } } func TestRestoreCleansMaintenanceAfterMutationAndRollbackFailures(t *testing.T) { installation := preflightTestInstallation(t) for _, test := range []struct { name string recoveryErr error }{ {name: "mutation"}, {name: "rollback", recoveryErr: errors.New("rollback failed")}, } { t.Run(test.name, func(t *testing.T) { backing := newBackupRunner(installation, true) deps := restoreTestDependencies(t, backing) mutationErr := errors.New("mutation failed") deps.restoreFile = func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error { return mutationErr } deps.recover = func(context.Context, config.Installation, PreflightResult, *stagedArchive, bool, authProjectionRestoreTransaction) error { if test.recoveryErr == nil { backing.running, backing.coreRunning = true, true } return test.recoveryErr } _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: restoreArchive(t), Confirm: true}, deps) if !errors.Is(err, mutationErr) || (test.recoveryErr != nil && !errors.Is(err, test.recoveryErr)) { t.Fatalf("restore error = %v, want mutation and rollback failures", err) } if test.recoveryErr == nil && (backing.maintenance || !backing.running) { t.Fatalf("successful recovery cleanup state: maintenance=%t running=%t", backing.maintenance, backing.running) } if test.recoveryErr != nil && !backing.maintenance { t.Fatal("failed recovery removed the maintenance barrier before a verified rollback") } }) } } func TestRestorePreservesDrainAndMaintenanceCleanupFailures(t *testing.T) { installation := preflightTestInstallation(t) backing := newBackupRunner(installation, true) backing.sessionResponses = []string{`[{"status":"running","archived":false}]`} cleanupErr := errors.New("maintenance cleanup failed") runner := &restoreFailureRunner{ fakeBackupRunner: backing, failContains: "operator-command.js maintenance-deactivate", err: cleanupErr, beforeFailure: func() { backing.maintenance = false }, } deps := restoreTestDependencies(t, runner) _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: restoreArchive(t), Confirm: true}, deps) if !errors.Is(err, ErrActiveSessions) || !errors.Is(err, cleanupErr) { t.Fatalf("restore error = %v, want drain and cleanup failures", err) } if backing.maintenance || !backing.running { t.Fatalf("cleanup failure state: maintenance=%t running=%t", backing.maintenance, backing.running) } } func restoreArchive(t *testing.T) string { t.Helper() archive := filepath.Join(t.TempDir(), "restore.zip") writePreflightArchive(t, archive, preflightArchiveSpec{ entries: []preflightArchiveEntry{{path: "configuration/operator.env", body: []byte("safe")}}, }) return archive } type failStartRestoreRunner struct { *fakeBackupRunner } func (runner failStartRestoreRunner) Run(ctx context.Context, args []string, stdin io.Reader) (compose.Result, error) { if strings.HasSuffix(strings.Join(args, " "), " start") { return compose.Result{}, errors.New("start refused") } return runner.fakeBackupRunner.Run(ctx, args, stdin) } type restoreFailureRunner struct { *fakeBackupRunner failContains string failSuffix string err error beforeFailure func() matchCount int } func (runner *restoreFailureRunner) Run(ctx context.Context, args []string, stdin io.Reader) (compose.Result, error) { command := strings.Join(args, " ") matches := runner.failContains != "" && strings.Contains(command, runner.failContains) matches = matches || runner.failSuffix != "" && strings.HasSuffix(command, runner.failSuffix) if matches { runner.matchCount++ if runner.beforeFailure != nil { runner.beforeFailure() } return compose.Result{}, runner.err } return runner.fakeBackupRunner.Run(ctx, args, stdin) } func (runner *restoreFailureRunner) 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 *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 authenticationStateResetRunner struct { args []string result compose.Result err error } func (runner *authenticationStateResetRunner) Run(_ context.Context, args []string, _ io.Reader) (compose.Result, error) { runner.args = append([]string(nil), args...) return runner.result, runner.err } func (runner *authenticationStateResetRunner) Stream(context.Context, []string, io.Reader, io.Writer) (compose.Result, error) { return compose.Result{}, errors.New("authentication state reset must not stream a volume archive") } func (*authenticationStateResetRunner) SessionInventoryScope() string { return "mine" } type workspaceVerificationRunner struct { running bool result compose.Result err error calls []string } func (runner *workspaceVerificationRunner) Run(_ context.Context, args []string, _ io.Reader) (compose.Result, error) { command := strings.Join(args, " ") runner.calls = append(runner.calls, command) if strings.Contains(command, " ps --all --format json") { if runner.running { return compose.Result{Stdout: healthyServicesPayload()}, nil } return compose.Result{}, nil } if strings.Contains(command, "operator-command.js workspace-integrity") { return runner.result, runner.err } return compose.Result{}, fmt.Errorf("unexpected workspace verification command: %s", command) } func (*workspaceVerificationRunner) Stream(context.Context, []string, io.Reader, io.Writer) (compose.Result, error) { return compose.Result{}, errors.New("workspace verification must not stream") } func (*workspaceVerificationRunner) SessionInventoryScope() string { return "mine" } func TestVerifyRestoreWorkspaceUsesFixedNonNetworkOperatorPath(t *testing.T) { installation := preflightTestInstallation(t) fingerprint := "sha256:" + strings.Repeat("a", 64) for _, test := range []struct { name string running bool payload string prefix string }{ {name: "stopped uninitialized", payload: fmt.Sprintf(`{"ready":true,"state":"uninitialized","workspaces":0,"fingerprint":%q}`, fingerprint), prefix: "run --rm --no-deps --no-TTY core"}, {name: "running active", running: true, payload: fmt.Sprintf(`{"ready":true,"state":"active","workspaces":1,"fingerprint":%q}`, fingerprint), prefix: "exec -T core"}, } { t.Run(test.name, func(t *testing.T) { runner := &workspaceVerificationRunner{ running: test.running, result: compose.Result{Stdout: test.payload}, } if err := verifyRestoreWorkspace(context.Background(), installation, runner); err != nil { t.Fatalf("verify restored workspace: %v", err) } if len(runner.calls) != 2 || !strings.Contains(runner.calls[1], test.prefix+" node /app/backend/dist/operator-command.js workspace-integrity") { t.Fatalf("workspace verification calls = %#v", runner.calls) } for _, call := range runner.calls { if strings.Contains(call, "curl") || strings.Contains(call, "-e const") || strings.Contains(strings.ToLower(call), "header") { t.Fatalf("workspace verification used an unsafe command: %s", call) } } }) } } func TestVerifyRestoreWorkspaceRejectsInvalidOperatorResults(t *testing.T) { installation := preflightTestInstallation(t) valid := `{"ready":true,"state":"active","workspaces":1,"fingerprint":"sha256:` + strings.Repeat("a", 64) + `"}` for _, test := range []struct { name string result compose.Result err error }{ {name: "empty"}, {name: "malformed", result: compose.Result{Stdout: `{malformed`}}, {name: "trailing document", result: compose.Result{Stdout: valid + `{}`}}, {name: "unknown field", result: compose.Result{Stdout: strings.TrimSuffix(valid, "}") + `,"detail":"unsafe"}`}}, {name: "not ready", result: compose.Result{Stdout: strings.Replace(valid, `"ready":true`, `"ready":false`, 1)}}, {name: "unknown state", result: compose.Result{Stdout: strings.Replace(valid, `"state":"active"`, `"state":"unknown"`, 1)}}, {name: "inconsistent count", result: compose.Result{Stdout: strings.Replace(valid, `"state":"active"`, `"state":"uninitialized"`, 1)}}, {name: "missing fingerprint", result: compose.Result{Stdout: `{"ready":true,"state":"active","workspaces":1}`}}, {name: "malformed fingerprint", result: compose.Result{Stdout: `{"ready":true,"state":"active","workspaces":1,"fingerprint":"sha256:not-a-digest"}`}}, {name: "nonzero", result: compose.Result{ExitCode: 2}, err: errors.New("exit status 2")}, } { t.Run(test.name, func(t *testing.T) { runner := &workspaceVerificationRunner{result: test.result, err: test.err} if err := verifyRestoreWorkspace(context.Background(), installation, runner); err == nil { t.Fatal("invalid workspace verification result was accepted") } }) } } func equalStrings(got, want []string) bool { if len(got) != len(want) { return false } for index := range got { if got[index] != want[index] { return false } } return true } func stageRecoveryForTest(t *testing.T, installation config.Installation, recovery PreflightResult) *stagedArchive { t.Helper() if err := os.MkdirAll(installation.ControlDirectory(), 0o700); err != nil { t.Fatal(err) } staged, err := recovery.StageArchive(context.Background()) if err != nil { t.Fatal(err) } t.Cleanup(func() { if err := staged.Close(); err != nil { t.Error(err) } }) return staged } func restoreTestDependencies(t *testing.T, runner archiveRunner) restoreDependencies { t.Helper() recoveryArchive := restoreArchive(t) return restoreDependencies{ preflight: func(ctx context.Context, installation config.Installation, request PreflightRequest) (PreflightResult, error) { return Preflight(ctx, installation, request, permissivePreflightDependencies()) }, checkpoint: func(context.Context, *lifecycle.Transaction, config.Installation, CreateRequest) (Result, error) { return Result{Path: "/tmp/default-checkpoint.zip"}, nil }, prepareRecovery: func(ctx context.Context, installation config.Installation, _ string) (PreflightResult, error) { return Preflight(ctx, installation, PreflightRequest{Archive: recoveryArchive, Confirm: true, AllowExternalSecrets: true}, permissivePreflightDependencies()) }, recover: func(context.Context, config.Installation, PreflightResult, *stagedArchive, bool, authProjectionRestoreTransaction) error { return nil }, cleanupCheckpoint: func(string) error { return nil }, acquireTransaction: lifecycle.AcquireTransaction, requireAuthProjection: func() error { return nil }, runner: runner, sleep: func(time.Duration) {}, restoreFile: func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error { return nil }, restoreVolume: func(context.Context, config.Installation, VolumeMetadata, io.Reader) error { return nil }, resetAuthenticationState: func(context.Context, config.Installation, archiveRunner) error { return nil }, verify: map[string]restoreVerify{ "health": func(context.Context, config.Installation, archiveRunner) error { return nil }, "doctor": func(context.Context, config.Installation, archiveRunner) error { return nil }, "pi": func(context.Context, config.Installation, archiveRunner) error { return nil }, "workspace": func(context.Context, config.Installation, archiveRunner) error { return nil }, }, } } type projectionRestoreStub struct { events *[]string blocked bool publishErr error restoreErr error closeErr error } func (stub *projectionRestoreStub) PublishCanonical() (authconfig.ProjectionStatus, error) { *stub.events = append(*stub.events, "publish") if stub.publishErr != nil { return authconfig.ProjectionStatus{}, stub.publishErr } stub.blocked = false return authconfig.ProjectionStatus{State: "ready", Generation: "g", CanonicalRevision: "sha256:g", Equal: true}, nil } func (stub *projectionRestoreStub) RestorePriorIfCanonicalUnchanged() error { *stub.events = append(*stub.events, "restore-prior") if stub.restoreErr != nil { return stub.restoreErr } stub.blocked = false return nil } func (stub *projectionRestoreStub) Close() error { *stub.events = append(*stub.events, "close") return stub.closeErr } type restartEventRunner struct { archiveRunner events *[]string } func (runner restartEventRunner) Run(ctx context.Context, args []string, input io.Reader) (compose.Result, error) { if strings.HasSuffix(strings.Join(args, " "), " start") { *runner.events = append(*runner.events, "restart") } return runner.archiveRunner.Run(ctx, args, input) } func (runner restartEventRunner) Stream(ctx context.Context, args []string, input io.Reader, output io.Writer) (compose.Result, error) { return runner.archiveRunner.Stream(ctx, args, input, output) } func (runner restartEventRunner) SessionInventoryScope() string { return runner.archiveRunner.SessionInventoryScope() } func projectedRestoreFixture(t *testing.T) (config.Installation, string) { t.Helper() installation := preflightTestInstallation(t) authRoot := filepath.Join(t.TempDir(), "canonical-auth") if err := os.Mkdir(authRoot, 0o700); err != nil { t.Fatal(err) } installation.Authentication.ConfigDirectory = authRoot installation.Authentication.RuntimeProjection = &config.RuntimeProjection{Directory: filepath.Join(t.TempDir(), "runtime-auth"), UID: 10001, GID: 10001} archive := filepath.Join(t.TempDir(), "auth-restore.zip") writePreflightArchive(t, archive, preflightArchiveSpec{includeSecrets: true, entries: []preflightArchiveEntry{{ path: "authentication-secrets/000-auth.yaml", body: []byte("candidate auth\n"), kind: EntryExternalSecret, sensitive: true, owner: "authentication-configuration", sourcePath: filepath.Join(authRoot, "auth.yaml"), }}}) return installation, archive } func assertOrderedEvents(t *testing.T, events []string, wants ...string) { t.Helper() at := 0 for _, want := range wants { for at < len(events) && events[at] != want { at++ } if at == len(events) { t.Fatalf("events = %v, want ordered subsequence %v", events, wants) } at++ } } func TestRestoreAuthBearingArchiveBlocksBeforeWritePublishesBeforeRestart(t *testing.T) { installation, archive := projectedRestoreFixture(t) var events []string backing := newBackupRunner(installation, true) runner := restartEventRunner{archiveRunner: backing, events: &events} deps := restoreTestDependencies(t, runner) transaction := &projectionRestoreStub{events: &events} deps.beginAuthProjection = func(context.Context, config.Installation) (authProjectionRestoreTransaction, error) { transaction.blocked = true events = append(events, "begin") return transaction, nil } deps.restoreFile = func(_ context.Context, _ config.Installation, entry ArchiveEntryMetadata, _ io.Reader) error { if entry.Owner == "authentication-configuration" { if !transaction.blocked { t.Fatal("auth destination write began without blocked projection") } events = append(events, "restore-auth") } return nil } if _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: archive, Confirm: true}, deps); err != nil { t.Fatal(err) } assertOrderedEvents(t, events, "begin", "restore-auth", "publish", "restart", "close") if transaction.blocked { t.Fatalf("successful restore closed while projection remained blocked: %v", events) } } func TestRestoreAuthCandidatePublicationFailureRecoversThenStaysBlockedWithoutRestart(t *testing.T) { installation, archive := projectedRestoreFixture(t) var events []string backing := newBackupRunner(installation, true) runner := restartEventRunner{archiveRunner: backing, events: &events} deps := restoreTestDependencies(t, runner) publicationErr := errors.New("synthetic publication failure") transaction := &projectionRestoreStub{events: &events, publishErr: publicationErr} deps.beginAuthProjection = func(context.Context, config.Installation) (authProjectionRestoreTransaction, error) { transaction.blocked = true return transaction, nil } deps.recover = func(_ context.Context, _ config.Installation, _ PreflightResult, _ *stagedArchive, _ bool, transaction authProjectionRestoreTransaction) error { events = append(events, "restore-checkpoint") _, err := transaction.PublishCanonical() return err } if _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: archive, Confirm: true}, deps); !errors.Is(err, publicationErr) { t.Fatalf("restore error = %v, want publication failure", err) } assertOrderedEvents(t, events, "publish", "restore-checkpoint", "close") if !transaction.blocked { t.Fatal("publication failure reopened projection") } if backing.startCount != 0 { t.Fatalf("restart after failed candidate/recovery publication = %d", backing.startCount) } if !backing.maintenance { t.Fatal("admissions reopened after failed recovery publication") } } func TestRestoreAuthVerificationFailuresRepublishCheckpointBeforeRecoveryRestart(t *testing.T) { for _, failed := range []string{"health", "doctor", "pi", "workspace"} { t.Run(failed, func(t *testing.T) { installation, archive := projectedRestoreFixture(t) var events []string backing := newBackupRunner(installation, true) runner := restartEventRunner{archiveRunner: backing, events: &events} deps := restoreTestDependencies(t, runner) transaction := &projectionRestoreStub{events: &events} deps.beginAuthProjection = func(context.Context, config.Installation) (authProjectionRestoreTransaction, error) { transaction.blocked = true return transaction, nil } for _, name := range []string{"health", "doctor", "pi", "workspace"} { name := name deps.verify[name] = func(context.Context, config.Installation, archiveRunner) error { events = append(events, "candidate-"+name) if name == failed { return errors.New("synthetic " + name + " failure") } return nil } } deps.recover = func(ctx context.Context, target config.Installation, _ PreflightResult, _ *stagedArchive, wasRunning bool, transaction authProjectionRestoreTransaction) error { events = append(events, "restore-checkpoint") if _, err := transaction.PublishCanonical(); err != nil { return err } events = append(events, "verify-checkpoint") if wasRunning { return composeStartAndVerify(ctx, target, runner) } return nil } if _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: archive, Confirm: true}, deps); err == nil { t.Fatalf("%s failure accepted", failed) } assertOrderedEvents(t, events, "publish", "candidate-"+failed, "restore-checkpoint", "publish", "verify-checkpoint", "restart", "close") if transaction.blocked { t.Fatalf("verified checkpoint recovery remained blocked: %v", events) } }) } } func TestRestoreAuthPreMutationFailureRestoresPriorReadySelector(t *testing.T) { installation, archive := projectedRestoreFixture(t) var events []string backing := newBackupRunner(installation, false) runner := &commandFailureRunner{fakeBackupRunner: backing, failures: []*commandFailure{{ match: func(command string) bool { return strings.Contains(command, " ps --all --format json") }, err: errors.New("synthetic pre-mutation failure"), remaining: 1, }}} deps := restoreTestDependencies(t, runner) transaction := &projectionRestoreStub{events: &events} deps.beginAuthProjection = func(context.Context, config.Installation) (authProjectionRestoreTransaction, error) { transaction.blocked = true return transaction, nil } deps.recover = func(context.Context, config.Installation, PreflightResult, *stagedArchive, bool, authProjectionRestoreTransaction) error { t.Fatal("recovery ran before any destination mutation") return nil } if _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: archive, Confirm: true}, deps); err == nil { t.Fatal("pre-mutation failure accepted") } assertOrderedEvents(t, events, "restore-prior", "close") if transaction.blocked { t.Fatal("unchanged canonical pre-write failure did not restore ready selector") } } func TestRestoreNonAuthArchiveNeverBeginsProjection(t *testing.T) { installation := preflightTestInstallation(t) deps := restoreTestDependencies(t, newBackupRunner(installation, false)) called := false deps.beginAuthProjection = func(context.Context, config.Installation) (authProjectionRestoreTransaction, error) { called = true return nil, errors.New("must not begin") } if _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: restoreArchive(t), Confirm: true}, deps); err != nil { t.Fatal(err) } if called { t.Fatal("non-auth archive began projection transaction") } }