2263 lines
89 KiB
Go
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")
|
|
}
|
|
}
|