diff --git a/tools/tht/internal/backup/restore.go b/tools/tht/internal/backup/restore.go new file mode 100644 index 00000000..ce045bbf --- /dev/null +++ b/tools/tht/internal/backup/restore.go @@ -0,0 +1,107 @@ +package backup + +import ( + "archive/zip" + "context" + "errors" + "fmt" + "io" + "time" + + "github.com/aritmolab/thothii/tools/tht/internal/config" +) + +var ErrRestoreConfirmationRequired = errors.New("restore requires --yes") + +// RestoreRequest describes the deliberately-confirmed archive restoration. +type RestoreRequest struct { + Archive string + Confirm bool + Drain bool +} + +// RestoreResult records the retained recovery point and the final service state. +type RestoreResult struct { + Checkpoint string + Restarted bool + Verified bool +} + +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) + acquireLock func(config.Installation) (restoreLock, error) + runner archiveRunner + sleep func(duration time.Duration) + restoreFile func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error + restoreVolume func(context.Context, config.Installation, VolumeMetadata, io.Reader) error + verify map[string]restoreVerify +} + +// Restore runs the host transaction. Concrete host dependencies are intentionally kept outside +// the deterministic core so callers cannot bypass its preflight and checkpoint boundaries. +func Restore(ctx context.Context, installation config.Installation, request RestoreRequest) (RestoreResult, error) { + return RestoreResult{}, errors.New("restore host dependencies are unavailable") +} + +func restoreWithDependencies(ctx context.Context, installation config.Installation, request RestoreRequest, deps restoreDependencies) (result RestoreResult, resultErr error) { + if !request.Confirm { return RestoreResult{}, ErrRestoreConfirmationRequired } + if request.Archive == "" { return RestoreResult{}, errors.New("restore archive is required") } + if deps.preflight == nil || deps.checkpoint == nil || deps.acquireLock == nil || deps.runner == nil || deps.restoreFile == nil || deps.restoreVolume == 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 { return RestoreResult{}, err } + + checkpoint, err := deps.checkpoint(ctx, installation, CreateRequest{}) + if err != nil { return RestoreResult{}, fmt.Errorf("create recovery checkpoint: %w", err) } + result.Checkpoint = checkpoint.Path + lock, err := deps.acquireLock(installation) + if err != nil { return result, err } + defer func() { if releaseErr := lock.Release(); releaseErr != nil && resultErr == nil { resultErr = releaseErr } }() + + wasRunning, err := installationRunning(ctx, installation, deps.runner) + if err != nil { return result, err } + mutated := false + defer func() { + if resultErr != nil && mutated { _ = runCompose(context.Background(), installation, deps.runner, "stop") } + }() + if wasRunning { + if err := maintenance(ctx, installation, deps.runner, true); err != nil { return result, err } + if err := waitForNoActiveSessions(ctx, installation, deps.runner, request.Drain, deps.sleep); err != nil { return result, err } + if err := runCompose(ctx, installation, deps.runner, "stop"); err != nil { return result, err } + } + reader, err := zip.NewReader(archive, preflight.ArchiveSize) + if err != nil { return result, fmt.Errorf("read verified restore archive: %w", err) } + members := make(map[string]*zip.File, len(reader.File)) + for _, member := range reader.File { members[member.Name] = member } + for _, entry := range preflight.Entries { + if entry.Kind == EntryVolume { continue } + member := members[entry.Path] + if member == nil { return result, fmt.Errorf("verified archive is missing %q", entry.Path) } + stream, openErr := member.Open() + if openErr != nil { return result, fmt.Errorf("open verified archive member %q: %w", entry.Path, openErr) } + mutated = true + restoreErr := deps.restoreFile(ctx, installation, entry, stream) + closeErr := stream.Close() + if restoreErr != nil { return result, restoreErr } + if closeErr != nil { return result, closeErr } + } + if wasRunning { + if err := composeStartAndVerify(ctx, installation, deps.runner); err != nil { return result, err } + result.Restarted = true + } + for _, name := range []string{"health", "doctor", "pi", "workspace"} { + check := deps.verify[name] + if check == nil { return result, fmt.Errorf("restore verification %q is unavailable", name) } + if err := check(ctx, installation, deps.runner); err != nil { return result, fmt.Errorf("restore verification %s: %w", name, err) } + } + result.Verified = true + return result, nil +} diff --git a/tools/tht/internal/backup/restore_test.go b/tools/tht/internal/backup/restore_test.go new file mode 100644 index 00000000..ee85b11d --- /dev/null +++ b/tools/tht/internal/backup/restore_test.go @@ -0,0 +1,258 @@ +package backup + +import ( + "context" + "errors" + "io" + "path/filepath" + "strings" + "testing" + "time" + + "github.com/aritmolab/thothii/tools/tht/internal/compose" + "github.com/aritmolab/thothii/tools/tht/internal/config" +) + +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) + deps.checkpoint = func(_ context.Context, _ config.Installation, request CreateRequest) (Result, error) { + events = append(events, "checkpoint") + checkpointRequest = request + return Result{Path: "/tmp/checkpoint.zip"}, nil + } + deps.acquireLock = func(config.Installation) (restoreLock, error) { + events = append(events, "lock") + return fakeRestoreLock{release: func() { events = append(events, "unlock") }}, nil + } + 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 != "/tmp/checkpoint.zip" || result.Restarted || !result.Verified { + t.Fatalf("Restore() result = %#v", result) + } + if checkpointRequest.IncludeSecrets || checkpointRequest.Confirm { + t.Fatalf("checkpoint request = %#v, want non-secret unconfirmed checkpoint", checkpointRequest) + } + if got, want := events, []string{"checkpoint", "lock", "file:configuration/operator.env", "health", "doctor", "pi", "workspace", "unlock"}; !equalStrings(got, want) { + t.Fatalf("restore events = %v, want %v", got, want) + } +} + +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, 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, 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 TestRestoreFileFailureStopsMutatedTargetAndRetainsCheckpoint(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, config.Installation, CreateRequest) (Result, error) { + return Result{Path: "/tmp/recovery.zip"}, nil + } + deps.restoreFile = func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error { + 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 != "/tmp/recovery.zip" { + t.Fatalf("recovery checkpoint = %q, want retained path", result.Checkpoint) + } + if runner.running || runner.stopCount != 2 { + t.Fatalf("mutated target was not stopped: running=%t stops=%d", runner.running, runner.stopCount) + } +} + +func TestRestoreStartFailureStopsRunningTarget(t *testing.T) { + installation := preflightTestInstallation(t) + archive := restoreArchive(t) + backingRunner := newBackupRunner(installation, true) + deps := restoreTestDependencies(t, failStartRestoreRunner{fakeBackupRunner: backingRunner}) + + _, 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 backingRunner.running || backingRunner.stopCount != 2 || backingRunner.startCount != 0 { + t.Fatalf("failed restart left target available: running=%t stops=%d starts=%d", backingRunner.running, backingRunner.stopCount, backingRunner.startCount) + } +} + +func TestRestoreVerificationFailureStopsRunningTarget(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 } + + _, 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 runner.running || runner.stopCount != 2 || runner.startCount != 1 { + t.Fatalf("verification failure left target available: running=%t stops=%d starts=%d", runner.running, runner.stopCount, runner.startCount) + } +} + +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 { + t.Fatalf("active-session refusal mutated target: files=%d stops=%d starts=%d running=%t", restoredFiles, runner.stopCount, runner.startCount, runner.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 fakeRestoreLock struct { + release func() +} + +func (lock fakeRestoreLock) Release() error { + if lock.release != nil { + lock.release() + } + return nil +} + +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 restoreTestDependencies(t *testing.T, runner archiveRunner) restoreDependencies { + t.Helper() + return restoreDependencies{ + 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) { + return Result{Path: "/tmp/default-checkpoint.zip"}, nil + }, + acquireLock: func(config.Installation) (restoreLock, error) { return fakeRestoreLock{}, 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 + }, + 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 }, + }, + } +}