feat(cli): add restore transaction core

This commit is contained in:
2026-08-16 02:05:21 +02:00
parent cbb18acc4d
commit 69dca0820b
2 changed files with 365 additions and 0 deletions
+107
View File
@@ -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
}
+258
View File
@@ -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 },
},
}
}