Files
ThothII/tools/tht/internal/backup/restore_test.go
T

2263 lines
89 KiB
Go

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,
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")
}
}