Files
marcopan 610ae8c85a fix(auth): make local verification portable
Keep upstream identity visible while limiting logout to local auth. Inject the restore privilege gate so the deterministic core tests do not depend on the host OS, and confine descriptor-backed projection tests to Linux. Accept the real remaining Pi timeout budget instead of an exact millisecond.
2026-08-25 10:50:30 +02:00

2264 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,
requireAuthProjection: func() error { return nil },
runner: runner,
sleep: func(time.Duration) {},
restoreFile: func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error { return nil },
restoreVolume: func(context.Context, config.Installation, VolumeMetadata, io.Reader) error {
return nil
},
resetAuthenticationState: func(context.Context, config.Installation, archiveRunner) error {
return nil
},
verify: map[string]restoreVerify{
"health": func(context.Context, config.Installation, archiveRunner) error { return nil },
"doctor": func(context.Context, config.Installation, archiveRunner) error { return nil },
"pi": func(context.Context, config.Installation, archiveRunner) error { return nil },
"workspace": func(context.Context, config.Installation, archiveRunner) error { return nil },
},
}
}
type projectionRestoreStub struct {
events *[]string
blocked bool
publishErr error
restoreErr error
closeErr error
}
func (stub *projectionRestoreStub) PublishCanonical() (authconfig.ProjectionStatus, error) {
*stub.events = append(*stub.events, "publish")
if stub.publishErr != nil {
return authconfig.ProjectionStatus{}, stub.publishErr
}
stub.blocked = false
return authconfig.ProjectionStatus{State: "ready", Generation: "g", CanonicalRevision: "sha256:g", Equal: true}, nil
}
func (stub *projectionRestoreStub) RestorePriorIfCanonicalUnchanged() error {
*stub.events = append(*stub.events, "restore-prior")
if stub.restoreErr != nil {
return stub.restoreErr
}
stub.blocked = false
return nil
}
func (stub *projectionRestoreStub) Close() error {
*stub.events = append(*stub.events, "close")
return stub.closeErr
}
type restartEventRunner struct {
archiveRunner
events *[]string
}
func (runner restartEventRunner) Run(ctx context.Context, args []string, input io.Reader) (compose.Result, error) {
if strings.HasSuffix(strings.Join(args, " "), " start") {
*runner.events = append(*runner.events, "restart")
}
return runner.archiveRunner.Run(ctx, args, input)
}
func (runner restartEventRunner) Stream(ctx context.Context, args []string, input io.Reader, output io.Writer) (compose.Result, error) {
return runner.archiveRunner.Stream(ctx, args, input, output)
}
func (runner restartEventRunner) SessionInventoryScope() string {
return runner.archiveRunner.SessionInventoryScope()
}
func projectedRestoreFixture(t *testing.T) (config.Installation, string) {
t.Helper()
installation := preflightTestInstallation(t)
authRoot := filepath.Join(t.TempDir(), "canonical-auth")
if err := os.Mkdir(authRoot, 0o700); err != nil {
t.Fatal(err)
}
installation.Authentication.ConfigDirectory = authRoot
installation.Authentication.RuntimeProjection = &config.RuntimeProjection{Directory: filepath.Join(t.TempDir(), "runtime-auth"), UID: 10001, GID: 10001}
archive := filepath.Join(t.TempDir(), "auth-restore.zip")
writePreflightArchive(t, archive, preflightArchiveSpec{includeSecrets: true, entries: []preflightArchiveEntry{{
path: "authentication-secrets/000-auth.yaml", body: []byte("candidate auth\n"), kind: EntryExternalSecret,
sensitive: true, owner: "authentication-configuration", sourcePath: filepath.Join(authRoot, "auth.yaml"),
}}})
return installation, archive
}
func assertOrderedEvents(t *testing.T, events []string, wants ...string) {
t.Helper()
at := 0
for _, want := range wants {
for at < len(events) && events[at] != want {
at++
}
if at == len(events) {
t.Fatalf("events = %v, want ordered subsequence %v", events, wants)
}
at++
}
}
func TestRestoreAuthBearingArchiveBlocksBeforeWritePublishesBeforeRestart(t *testing.T) {
installation, archive := projectedRestoreFixture(t)
var events []string
backing := newBackupRunner(installation, true)
runner := restartEventRunner{archiveRunner: backing, events: &events}
deps := restoreTestDependencies(t, runner)
transaction := &projectionRestoreStub{events: &events}
deps.beginAuthProjection = func(context.Context, config.Installation) (authProjectionRestoreTransaction, error) {
transaction.blocked = true
events = append(events, "begin")
return transaction, nil
}
deps.restoreFile = func(_ context.Context, _ config.Installation, entry ArchiveEntryMetadata, _ io.Reader) error {
if entry.Owner == "authentication-configuration" {
if !transaction.blocked {
t.Fatal("auth destination write began without blocked projection")
}
events = append(events, "restore-auth")
}
return nil
}
if _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: archive, Confirm: true}, deps); err != nil {
t.Fatal(err)
}
assertOrderedEvents(t, events, "begin", "restore-auth", "publish", "restart", "close")
if transaction.blocked {
t.Fatalf("successful restore closed while projection remained blocked: %v", events)
}
}
func TestRestoreAuthCandidatePublicationFailureRecoversThenStaysBlockedWithoutRestart(t *testing.T) {
installation, archive := projectedRestoreFixture(t)
var events []string
backing := newBackupRunner(installation, true)
runner := restartEventRunner{archiveRunner: backing, events: &events}
deps := restoreTestDependencies(t, runner)
publicationErr := errors.New("synthetic publication failure")
transaction := &projectionRestoreStub{events: &events, publishErr: publicationErr}
deps.beginAuthProjection = func(context.Context, config.Installation) (authProjectionRestoreTransaction, error) {
transaction.blocked = true
return transaction, nil
}
deps.recover = func(_ context.Context, _ config.Installation, _ PreflightResult, _ *stagedArchive, _ bool, transaction authProjectionRestoreTransaction) error {
events = append(events, "restore-checkpoint")
_, err := transaction.PublishCanonical()
return err
}
if _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: archive, Confirm: true}, deps); !errors.Is(err, publicationErr) {
t.Fatalf("restore error = %v, want publication failure", err)
}
assertOrderedEvents(t, events, "publish", "restore-checkpoint", "close")
if !transaction.blocked {
t.Fatal("publication failure reopened projection")
}
if backing.startCount != 0 {
t.Fatalf("restart after failed candidate/recovery publication = %d", backing.startCount)
}
if !backing.maintenance {
t.Fatal("admissions reopened after failed recovery publication")
}
}
func TestRestoreAuthVerificationFailuresRepublishCheckpointBeforeRecoveryRestart(t *testing.T) {
for _, failed := range []string{"health", "doctor", "pi", "workspace"} {
t.Run(failed, func(t *testing.T) {
installation, archive := projectedRestoreFixture(t)
var events []string
backing := newBackupRunner(installation, true)
runner := restartEventRunner{archiveRunner: backing, events: &events}
deps := restoreTestDependencies(t, runner)
transaction := &projectionRestoreStub{events: &events}
deps.beginAuthProjection = func(context.Context, config.Installation) (authProjectionRestoreTransaction, error) {
transaction.blocked = true
return transaction, nil
}
for _, name := range []string{"health", "doctor", "pi", "workspace"} {
name := name
deps.verify[name] = func(context.Context, config.Installation, archiveRunner) error {
events = append(events, "candidate-"+name)
if name == failed {
return errors.New("synthetic " + name + " failure")
}
return nil
}
}
deps.recover = func(ctx context.Context, target config.Installation, _ PreflightResult, _ *stagedArchive, wasRunning bool, transaction authProjectionRestoreTransaction) error {
events = append(events, "restore-checkpoint")
if _, err := transaction.PublishCanonical(); err != nil {
return err
}
events = append(events, "verify-checkpoint")
if wasRunning {
return composeStartAndVerify(ctx, target, runner)
}
return nil
}
if _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: archive, Confirm: true}, deps); err == nil {
t.Fatalf("%s failure accepted", failed)
}
assertOrderedEvents(t, events, "publish", "candidate-"+failed, "restore-checkpoint", "publish", "verify-checkpoint", "restart", "close")
if transaction.blocked {
t.Fatalf("verified checkpoint recovery remained blocked: %v", events)
}
})
}
}
func TestRestoreAuthPreMutationFailureRestoresPriorReadySelector(t *testing.T) {
installation, archive := projectedRestoreFixture(t)
var events []string
backing := newBackupRunner(installation, false)
runner := &commandFailureRunner{fakeBackupRunner: backing, failures: []*commandFailure{{
match: func(command string) bool { return strings.Contains(command, " ps --all --format json") },
err: errors.New("synthetic pre-mutation failure"), remaining: 1,
}}}
deps := restoreTestDependencies(t, runner)
transaction := &projectionRestoreStub{events: &events}
deps.beginAuthProjection = func(context.Context, config.Installation) (authProjectionRestoreTransaction, error) {
transaction.blocked = true
return transaction, nil
}
deps.recover = func(context.Context, config.Installation, PreflightResult, *stagedArchive, bool, authProjectionRestoreTransaction) error {
t.Fatal("recovery ran before any destination mutation")
return nil
}
if _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: archive, Confirm: true}, deps); err == nil {
t.Fatal("pre-mutation failure accepted")
}
assertOrderedEvents(t, events, "restore-prior", "close")
if transaction.blocked {
t.Fatal("unchanged canonical pre-write failure did not restore ready selector")
}
}
func TestRestoreNonAuthArchiveNeverBeginsProjection(t *testing.T) {
installation := preflightTestInstallation(t)
deps := restoreTestDependencies(t, newBackupRunner(installation, false))
called := false
deps.beginAuthProjection = func(context.Context, config.Installation) (authProjectionRestoreTransaction, error) {
called = true
return nil, errors.New("must not begin")
}
if _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: restoreArchive(t), Confirm: true}, deps); err != nil {
t.Fatal(err)
}
if called {
t.Fatal("non-auth archive began projection transaction")
}
}