test(auth): gate local and OIDC authentication release
This commit is contained in:
@@ -127,23 +127,32 @@ func createWithDependencies(ctx context.Context, installation config.Installatio
|
||||
if request.IncludeSecrets && !request.Confirm {
|
||||
return Result{}, ErrConfirmationRequired
|
||||
}
|
||||
lock, err := lifecycle.Acquire(installation)
|
||||
transaction, err := lifecycle.AcquireTransaction(installation)
|
||||
if err != nil {
|
||||
return Result{}, err
|
||||
}
|
||||
defer func() {
|
||||
if releaseErr := lock.Release(); releaseErr != nil {
|
||||
if releaseErr := transaction.Release(); releaseErr != nil {
|
||||
result = Result{}
|
||||
resultErr = errors.Join(resultErr, fmt.Errorf("release backup lifecycle lock: %w", releaseErr))
|
||||
}
|
||||
}()
|
||||
return createWithDependenciesLockHeld(ctx, installation, request, dependencies)
|
||||
return createWithDependenciesTransaction(ctx, transaction, installation, request, dependencies)
|
||||
}
|
||||
|
||||
// createWithDependenciesLockHeld performs backup creation while the caller owns the installation
|
||||
// lifecycle lock. It must never acquire a lifecycle lock itself: Restore uses this primitive to
|
||||
// create its recovery checkpoint inside its already-locked transaction.
|
||||
func createWithDependenciesLockHeld(ctx context.Context, installation config.Installation, request CreateRequest, dependencies dependencies) (result Result, resultErr error) {
|
||||
// createWithDependenciesTransaction performs backup creation only with an active, opaque
|
||||
// installation-bound lifecycle capability. Restore passes the capability it acquired for the
|
||||
// enclosing transaction, so a recovery checkpoint cannot run without the same lifecycle lock.
|
||||
func createWithDependenciesTransaction(ctx context.Context, transaction *lifecycle.Transaction, installation config.Installation, request CreateRequest, dependencies dependencies) (result Result, resultErr error) {
|
||||
if err := transaction.Verify(installation); err != nil {
|
||||
return Result{}, err
|
||||
}
|
||||
if dependencies.runner == nil || dependencies.now == nil || dependencies.homeDir == nil || dependencies.revision == nil || dependencies.sleep == nil || dependencies.reserveOutput == nil || dependencies.publishReserved == nil {
|
||||
return Result{}, errors.New("backup dependencies are incomplete")
|
||||
}
|
||||
if request.IncludeSecrets && !request.Confirm {
|
||||
return Result{}, ErrConfirmationRequired
|
||||
}
|
||||
if err := rejectInlineSecretValues(installation.EnvFile); err != nil {
|
||||
return Result{}, err
|
||||
}
|
||||
|
||||
@@ -537,6 +537,33 @@ func TestCreateHonorsTheSharedInstallationLifecycleLock(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestCreateTransactionCapabilityRefusesUnlockedOrForeignInstallation(t *testing.T) {
|
||||
fixture := newBackupFixture(t, "local")
|
||||
runner := newBackupRunner(fixture.installation, false)
|
||||
request := CreateRequest{Output: filepath.Join(t.TempDir(), "transaction.zip")}
|
||||
|
||||
if _, err := createWithDependenciesTransaction(context.Background(), nil, fixture.installation, request, testDependencies(t, runner)); !errors.Is(err, lifecycle.ErrTransactionInactive) {
|
||||
t.Fatalf("nil transaction error = %v, want ErrTransactionInactive", err)
|
||||
}
|
||||
if len(runner.calls) != 0 {
|
||||
t.Fatalf("Docker runner was called without a transaction capability: %v", runner.calls)
|
||||
}
|
||||
|
||||
otherRoot := t.TempDir()
|
||||
other := config.Installation{ProjectDirectory: otherRoot, Path: filepath.Join(otherRoot, "thothii-installation.yaml")}
|
||||
transaction, err := lifecycle.AcquireTransaction(other)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(func() { _ = transaction.Release() })
|
||||
if _, err := createWithDependenciesTransaction(context.Background(), transaction, fixture.installation, request, testDependencies(t, runner)); !errors.Is(err, lifecycle.ErrTransactionInstallation) {
|
||||
t.Fatalf("foreign transaction error = %v, want ErrTransactionInstallation", err)
|
||||
}
|
||||
if len(runner.calls) != 0 {
|
||||
t.Fatalf("Docker runner was called with a foreign transaction capability: %v", runner.calls)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCreateRefusesMutableNonRunningServiceStates(t *testing.T) {
|
||||
for _, state := range []string{"paused", "restarting", "created", "removing"} {
|
||||
t.Run(state, func(t *testing.T) {
|
||||
|
||||
@@ -91,6 +91,14 @@ type verifiedArchive struct {
|
||||
limits PreflightLimits
|
||||
}
|
||||
|
||||
// stagedArchive holds an installation-private, immutable copy of the exact bytes accepted by
|
||||
// Preflight. The original archive remains retained only for provenance revalidation.
|
||||
type stagedArchive struct {
|
||||
file *os.File
|
||||
path string
|
||||
directory string
|
||||
}
|
||||
|
||||
type inspectedArchiveEntry struct {
|
||||
metadata ArchiveEntryMetadata
|
||||
member *zip.File
|
||||
@@ -176,12 +184,17 @@ func Preflight(ctx context.Context, installation config.Installation, request Pr
|
||||
if err != nil {
|
||||
return PreflightResult{}, err
|
||||
}
|
||||
stagingBytes := uint64(openedInfo.Size())
|
||||
if stagingBytes > ^uint64(0)-requiredBytes {
|
||||
return PreflightResult{}, errors.New("restore staging requirement exceeds supported size")
|
||||
}
|
||||
requiredWithStaging := requiredBytes + stagingBytes
|
||||
freeBytes, err := dependencies.FreeBytes(installation.ProjectDirectory)
|
||||
if err != nil {
|
||||
return PreflightResult{}, fmt.Errorf("check free disk space: %w", err)
|
||||
}
|
||||
if freeBytes < requiredBytes {
|
||||
return PreflightResult{}, fmt.Errorf("insufficient free disk space for restore: need %d bytes, have %d", requiredBytes, freeBytes)
|
||||
if freeBytes < requiredWithStaging {
|
||||
return PreflightResult{}, fmt.Errorf("insufficient free disk space for restore: need %d bytes, have %d", requiredWithStaging, freeBytes)
|
||||
}
|
||||
if err := contextError(ctx); err != nil {
|
||||
return PreflightResult{}, err
|
||||
@@ -214,6 +227,10 @@ func Preflight(ctx context.Context, installation config.Installation, request Pr
|
||||
// it immediately before a restore transaction and use the returned retained handle, never reopen
|
||||
// ArchivePath. It refuses a path replacement or in-place content change.
|
||||
func (result PreflightResult) RevalidateArchive() (*os.File, error) {
|
||||
return result.revalidateArchive(context.Background())
|
||||
}
|
||||
|
||||
func (result PreflightResult) revalidateArchive(ctx context.Context) (*os.File, error) {
|
||||
if result.archive == nil || result.archive.file == nil {
|
||||
return nil, errors.New("backup archive has not been retained by preflight")
|
||||
}
|
||||
@@ -228,13 +245,115 @@ func (result PreflightResult) RevalidateArchive() (*os.File, error) {
|
||||
if err != nil || !os.SameFile(result.archive.info, heldInfo) || heldInfo.Size() != result.ArchiveSize {
|
||||
return nil, errors.New("backup archive changed after preflight")
|
||||
}
|
||||
digest, err := digestArchive(context.Background(), result.archive.file, result.archive.limits.MaxArchiveBytes)
|
||||
digest, err := digestArchive(ctx, result.archive.file, result.archive.limits.MaxArchiveBytes)
|
||||
if err != nil || digest != result.archive.digest {
|
||||
return nil, errors.New("backup archive changed after preflight")
|
||||
}
|
||||
return result.archive.file, nil
|
||||
}
|
||||
|
||||
// StageArchive revalidates the retained archive and copies its exact bytes into a private file
|
||||
// immediately before extraction. Later writes to the source archive cannot affect extraction.
|
||||
func (result PreflightResult) StageArchive(ctx context.Context) (_ *stagedArchive, resultErr error) {
|
||||
source, err := result.revalidateArchive(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
directory, err := os.MkdirTemp("", "tht-restore-stage-")
|
||||
if err != nil {
|
||||
return nil, errors.New("create private restore staging directory")
|
||||
}
|
||||
if err := os.Chmod(directory, 0o700); err != nil {
|
||||
_ = os.Remove(directory)
|
||||
return nil, errors.New("protect private restore staging directory")
|
||||
}
|
||||
path := filepath.Join(directory, "archive.zip")
|
||||
file, err := os.OpenFile(path, os.O_RDWR|os.O_CREATE|os.O_EXCL, 0o600)
|
||||
if err != nil {
|
||||
_ = os.Remove(directory)
|
||||
return nil, errors.New("create private restore staging archive")
|
||||
}
|
||||
staged := &stagedArchive{file: file, path: path, directory: directory}
|
||||
completed := false
|
||||
defer func() {
|
||||
if !completed {
|
||||
_ = staged.Close()
|
||||
}
|
||||
}()
|
||||
if _, err := source.Seek(0, io.SeekStart); err != nil {
|
||||
return nil, errors.New("seek verified backup archive for staging")
|
||||
}
|
||||
digest := sha256.New()
|
||||
buffer := make([]byte, 128*1024)
|
||||
var total int64
|
||||
for {
|
||||
if err := contextError(ctx); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
count, readErr := source.Read(buffer)
|
||||
if count > 0 {
|
||||
if int64(count) > result.ArchiveSize-total {
|
||||
return nil, errors.New("backup archive changed after preflight")
|
||||
}
|
||||
written, writeErr := file.Write(buffer[:count])
|
||||
if writeErr != nil || written != count {
|
||||
return nil, errors.New("write private restore staging archive")
|
||||
}
|
||||
if _, writeErr := digest.Write(buffer[:count]); writeErr != nil {
|
||||
return nil, errors.New("hash private restore staging archive")
|
||||
}
|
||||
total += int64(count)
|
||||
}
|
||||
if errors.Is(readErr, io.EOF) {
|
||||
break
|
||||
}
|
||||
if readErr != nil {
|
||||
return nil, errors.New("read verified backup archive for staging")
|
||||
}
|
||||
}
|
||||
if total != result.ArchiveSize || digestForHash(digest) != result.archive.digest {
|
||||
return nil, errors.New("backup archive changed after preflight")
|
||||
}
|
||||
if err := file.Sync(); err != nil {
|
||||
return nil, errors.New("sync private restore staging archive")
|
||||
}
|
||||
if _, err := file.Seek(0, io.SeekStart); err != nil {
|
||||
return nil, errors.New("rewind private restore staging archive")
|
||||
}
|
||||
completed = true
|
||||
return staged, nil
|
||||
}
|
||||
|
||||
// Close removes only the staging file and directory created by StageArchive.
|
||||
func (staged *stagedArchive) Close() error {
|
||||
if staged == nil {
|
||||
return nil
|
||||
}
|
||||
var failed bool
|
||||
if staged.file != nil {
|
||||
if err := staged.file.Close(); err != nil {
|
||||
failed = true
|
||||
}
|
||||
staged.file = nil
|
||||
}
|
||||
if staged.path != "" {
|
||||
if err := os.Remove(staged.path); err != nil && !errors.Is(err, os.ErrNotExist) {
|
||||
failed = true
|
||||
}
|
||||
staged.path = ""
|
||||
}
|
||||
if staged.directory != "" {
|
||||
if err := os.Remove(staged.directory); err != nil && !errors.Is(err, os.ErrNotExist) {
|
||||
failed = true
|
||||
}
|
||||
staged.directory = ""
|
||||
}
|
||||
if failed {
|
||||
return errors.New("destroy private restore staging archive")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// CloseArchive releases the retained read-only archive handle after the caller finishes the
|
||||
// restore transaction or decides not to proceed.
|
||||
func (result PreflightResult) CloseArchive() error {
|
||||
|
||||
@@ -7,6 +7,7 @@ import (
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"errors"
|
||||
"io"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
@@ -356,6 +357,79 @@ func TestPreflightRevalidationRefusesAnArchivePathThatWasReplaced(t *testing.T)
|
||||
}
|
||||
}
|
||||
|
||||
func TestPreflightStagesArchiveIntoImmutablePrivateBytes(t *testing.T) {
|
||||
installation := preflightTestInstallation(t)
|
||||
archive := filepath.Join(t.TempDir(), "checked.zip")
|
||||
writePreflightArchive(t, archive, preflightArchiveSpec{
|
||||
entries: []preflightArchiveEntry{{path: "configuration/operator.env", body: []byte("before")}},
|
||||
})
|
||||
|
||||
result, err := Preflight(context.Background(), installation, PreflightRequest{Archive: archive, Confirm: true}, permissivePreflightDependencies())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer result.CloseArchive()
|
||||
staged, err := result.StageArchive(context.Background())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer staged.Close()
|
||||
|
||||
writePreflightArchive(t, archive, preflightArchiveSpec{
|
||||
entries: []preflightArchiveEntry{{path: "configuration/operator.env", body: []byte("after!")}},
|
||||
})
|
||||
reader, err := zip.NewReader(staged.file, result.ArchiveSize)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(reader.File) < 2 {
|
||||
t.Fatalf("staged archive members = %d, want manifest and payload", len(reader.File))
|
||||
}
|
||||
var payload *zip.File
|
||||
for _, member := range reader.File {
|
||||
if member.Name == "configuration/operator.env" {
|
||||
payload = member
|
||||
break
|
||||
}
|
||||
}
|
||||
if payload == nil {
|
||||
t.Fatal("staged archive is missing the configured payload")
|
||||
}
|
||||
stream, err := payload.Open()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer stream.Close()
|
||||
body, err := io.ReadAll(stream)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if string(body) != "before" {
|
||||
t.Fatalf("staged payload = %q, want preflighted bytes", body)
|
||||
}
|
||||
}
|
||||
|
||||
func TestPreflightStagingRejectsInPlaceArchiveHashMutation(t *testing.T) {
|
||||
installation := preflightTestInstallation(t)
|
||||
archive := filepath.Join(t.TempDir(), "checked.zip")
|
||||
writePreflightArchive(t, archive, preflightArchiveSpec{
|
||||
entries: []preflightArchiveEntry{{path: "configuration/operator.env", body: []byte("before")}},
|
||||
})
|
||||
|
||||
result, err := Preflight(context.Background(), installation, PreflightRequest{Archive: archive, Confirm: true}, permissivePreflightDependencies())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer result.CloseArchive()
|
||||
writePreflightArchive(t, archive, preflightArchiveSpec{
|
||||
entries: []preflightArchiveEntry{{path: "configuration/operator.env", body: []byte("after!")}},
|
||||
})
|
||||
|
||||
if _, err := result.StageArchive(context.Background()); err == nil || !strings.Contains(err.Error(), "changed") {
|
||||
t.Fatalf("StageArchive() error = %v, want changed archive refusal", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestPreflightRejectsArchiveEntryWithModeDifferentFromManifest(t *testing.T) {
|
||||
installation := preflightTestInstallation(t)
|
||||
archive := filepath.Join(t.TempDir(), "mode-mismatch.zip")
|
||||
|
||||
@@ -9,6 +9,7 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/aritmolab/thothii/tools/tht/internal/config"
|
||||
"github.com/aritmolab/thothii/tools/tht/internal/lifecycle"
|
||||
)
|
||||
|
||||
var ErrRestoreConfirmationRequired = errors.New("restore requires --yes")
|
||||
@@ -29,19 +30,17 @@ type RestoreResult struct {
|
||||
Verified bool
|
||||
}
|
||||
|
||||
type restoreLock interface{ Release() error }
|
||||
type restoreVerify func(context.Context, config.Installation, archiveRunner) error
|
||||
|
||||
type restoreDependencies struct {
|
||||
preflight func(context.Context, config.Installation, PreflightRequest) (PreflightResult, error)
|
||||
// checkpointLocked creates the secret-aware recovery archive while the caller already owns
|
||||
// the installation lifecycle lock. It must not call public Create, which would re-acquire the
|
||||
// non-reentrant lock and deadlock the restore transaction.
|
||||
checkpointLocked func(context.Context, config.Installation, CreateRequest) (Result, error)
|
||||
// checkpoint requires the opaque capability created by lifecycle acquisition. It must not call
|
||||
// public Create, which would re-acquire the non-reentrant lock and deadlock the transaction.
|
||||
checkpoint func(context.Context, *lifecycle.Transaction, config.Installation, CreateRequest) (Result, error)
|
||||
prepareRecovery func(context.Context, config.Installation, string) (PreflightResult, error)
|
||||
recover func(context.Context, config.Installation, PreflightResult, bool) error
|
||||
cleanupCheckpoint func(string) error
|
||||
acquireLock func(config.Installation) (restoreLock, error)
|
||||
acquireTransaction func(config.Installation) (*lifecycle.Transaction, error)
|
||||
runner archiveRunner
|
||||
sleep func(duration time.Duration)
|
||||
restoreFile func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error
|
||||
@@ -84,36 +83,32 @@ func restoreWithDependencies(ctx context.Context, installation config.Installati
|
||||
if request.Archive == "" {
|
||||
return RestoreResult{}, errors.New("restore archive is required")
|
||||
}
|
||||
if deps.preflight == nil || deps.checkpointLocked == nil || deps.prepareRecovery == nil || deps.recover == nil || deps.cleanupCheckpoint == nil || deps.acquireLock == nil || deps.runner == nil || deps.restoreFile == nil || deps.restoreVolume == nil || deps.resetAuthenticationState == nil || deps.verify == nil {
|
||||
if deps.preflight == nil || deps.checkpoint == nil || deps.prepareRecovery == nil || deps.recover == nil || deps.cleanupCheckpoint == nil || deps.acquireTransaction == nil || deps.runner == nil || deps.restoreFile == nil || deps.restoreVolume == nil || deps.resetAuthenticationState == nil || deps.verify == nil {
|
||||
return RestoreResult{}, errors.New("restore dependencies are incomplete")
|
||||
}
|
||||
preflight, err := deps.preflight(ctx, installation, PreflightRequest{Archive: request.Archive, Confirm: true, AllowExternalSecrets: true})
|
||||
if err != nil {
|
||||
return RestoreResult{}, err
|
||||
}
|
||||
archive, err := preflight.RevalidateArchive()
|
||||
if err != nil {
|
||||
_ = preflight.CloseArchive()
|
||||
return RestoreResult{}, err
|
||||
}
|
||||
|
||||
// Lock ordering is lifecycle lock -> Compose/operator maintenance barrier. The core operator
|
||||
// command never acquires the host lifecycle lock, so this order cannot form a lock cycle with
|
||||
// Docker Compose or the durable maintenance marker.
|
||||
lock, err := deps.acquireLock(installation)
|
||||
transaction, err := deps.acquireTransaction(installation)
|
||||
if err != nil {
|
||||
_ = preflight.CloseArchive()
|
||||
return result, err
|
||||
}
|
||||
defer func() {
|
||||
if releaseErr := lock.Release(); releaseErr != nil {
|
||||
if releaseErr := transaction.Release(); releaseErr != nil {
|
||||
result = RestoreResult{}
|
||||
resultErr = errors.Join(resultErr, fmt.Errorf("release restore lifecycle lock: %w", releaseErr))
|
||||
}
|
||||
}()
|
||||
// Target-dependent preflight is intentionally inside the lifecycle transaction. This binds
|
||||
// target ownership, volume, image, and free-space checks to the later mutation.
|
||||
preflight, err := deps.preflight(ctx, installation, PreflightRequest{Archive: request.Archive, Confirm: true, AllowExternalSecrets: true})
|
||||
if err != nil {
|
||||
return RestoreResult{}, err
|
||||
}
|
||||
defer preflight.CloseArchive()
|
||||
|
||||
checkpoint, err := deps.checkpointLocked(ctx, installation, CreateRequest{IncludeSecrets: true, Confirm: true})
|
||||
checkpoint, err := deps.checkpoint(ctx, transaction, installation, CreateRequest{IncludeSecrets: true, Confirm: true})
|
||||
if err != nil {
|
||||
return result, fmt.Errorf("create recovery checkpoint: %w", err)
|
||||
}
|
||||
@@ -217,8 +212,18 @@ func restoreWithDependencies(ctx context.Context, installation config.Installati
|
||||
return result, err
|
||||
}
|
||||
}
|
||||
staged, err := preflight.StageArchive(ctx)
|
||||
if err != nil {
|
||||
return result, err
|
||||
}
|
||||
defer func() {
|
||||
if closeErr := staged.Close(); closeErr != nil {
|
||||
result = RestoreResult{}
|
||||
resultErr = errors.Join(resultErr, closeErr)
|
||||
}
|
||||
}()
|
||||
state.mutated = true
|
||||
if err := restoreVerifiedEntries(ctx, installation, preflight, archive, deps.restoreFile, deps.restoreVolume); err != nil {
|
||||
if err := restoreVerifiedEntries(ctx, installation, preflight, staged.file, deps.restoreFile, deps.restoreVolume); err != nil {
|
||||
return result, err
|
||||
}
|
||||
if err := deps.resetAuthenticationState(ctx, installation, deps.runner); err != nil {
|
||||
|
||||
@@ -41,21 +41,19 @@ func productionRestoreDependencies(installation config.Installation) restoreDepe
|
||||
},
|
||||
})
|
||||
},
|
||||
checkpointLocked: func(ctx context.Context, target config.Installation, request CreateRequest) (Result, error) {
|
||||
checkpoint: func(ctx context.Context, transaction *lifecycle.Transaction, target config.Installation, request CreateRequest) (Result, error) {
|
||||
path, err := restoreCheckpointPath(target, time.Now().UTC())
|
||||
if err != nil {
|
||||
return Result{}, err
|
||||
}
|
||||
request.Output = path
|
||||
return createWithDependenciesLockHeld(ctx, target, request, productionCreateDependencies(target))
|
||||
return createWithDependenciesTransaction(ctx, transaction, target, request, productionCreateDependencies(target))
|
||||
},
|
||||
cleanupCheckpoint: cleanupRecoveryCheckpoint,
|
||||
acquireLock: func(target config.Installation) (restoreLock, error) {
|
||||
return lifecycle.Acquire(target)
|
||||
},
|
||||
runner: runner,
|
||||
sleep: time.Sleep,
|
||||
restoreFile: restoreFilePayload,
|
||||
cleanupCheckpoint: cleanupRecoveryCheckpoint,
|
||||
acquireTransaction: lifecycle.AcquireTransaction,
|
||||
runner: runner,
|
||||
sleep: time.Sleep,
|
||||
restoreFile: restoreFilePayload,
|
||||
restoreVolume: func(ctx context.Context, _ config.Installation, volume VolumeMetadata, input io.Reader) error {
|
||||
result, err := runner.Stream(ctx, volumeRestoreCommand(volume.Name), input, io.Discard)
|
||||
if err != nil || result.ExitCode != 0 {
|
||||
@@ -90,8 +88,7 @@ func cleanupRecoveryCheckpoint(path string) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func recoverRestoreTransaction(ctx context.Context, installation config.Installation, recovery PreflightResult, wasRunning bool, deps restoreDependencies) error {
|
||||
var resultErr error
|
||||
func recoverRestoreTransaction(ctx context.Context, installation config.Installation, recovery PreflightResult, wasRunning bool, deps restoreDependencies) (resultErr error) {
|
||||
if err := runCompose(ctx, installation, deps.runner, "stop"); err != nil {
|
||||
resultErr = errors.Join(resultErr, err)
|
||||
// The first stop may already have taken effect before Docker lost its response. Retry the
|
||||
@@ -100,11 +97,16 @@ func recoverRestoreTransaction(ctx context.Context, installation config.Installa
|
||||
return errors.Join(resultErr, retryErr)
|
||||
}
|
||||
}
|
||||
archive, err := recovery.RevalidateArchive()
|
||||
staged, err := recovery.StageArchive(ctx)
|
||||
if err != nil {
|
||||
return errors.Join(resultErr, err)
|
||||
}
|
||||
if err := restoreVerifiedEntries(ctx, installation, recovery, archive, deps.restoreFile, deps.restoreVolume); err != nil {
|
||||
defer func() {
|
||||
if closeErr := staged.Close(); closeErr != nil {
|
||||
resultErr = errors.Join(resultErr, closeErr)
|
||||
}
|
||||
}()
|
||||
if err := restoreVerifiedEntries(ctx, installation, recovery, staged.file, deps.restoreFile, deps.restoreVolume); err != nil {
|
||||
return errors.Join(resultErr, err)
|
||||
}
|
||||
// Recovery restores only configuration/secret files and durable application volumes. Runtime
|
||||
|
||||
@@ -120,7 +120,7 @@ func TestRestoreStoppedInstallationRunsCheckpointRestoreAndVerification(t *testi
|
||||
var events []string
|
||||
var checkpointRequest CreateRequest
|
||||
deps := restoreTestDependencies(t, runner)
|
||||
deps.checkpointLocked = func(_ context.Context, _ config.Installation, request CreateRequest) (Result, error) {
|
||||
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
|
||||
@@ -133,9 +133,9 @@ func TestRestoreStoppedInstallationRunsCheckpointRestoreAndVerification(t *testi
|
||||
events = append(events, "cleanup-checkpoint")
|
||||
return nil
|
||||
}
|
||||
deps.acquireLock = func(config.Installation) (restoreLock, error) {
|
||||
deps.acquireTransaction = func(target config.Installation) (*lifecycle.Transaction, error) {
|
||||
events = append(events, "lock")
|
||||
return fakeRestoreLock{release: func() { events = append(events, "unlock") }}, nil
|
||||
return lifecycle.AcquireTransaction(target)
|
||||
}
|
||||
deps.restoreFile = func(_ context.Context, _ config.Installation, entry ArchiveEntryMetadata, _ io.Reader) error {
|
||||
events = append(events, "file:"+entry.Path)
|
||||
@@ -159,9 +159,12 @@ func TestRestoreStoppedInstallationRunsCheckpointRestoreAndVerification(t *testi
|
||||
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", "unlock"}; !equalStrings(got, want) {
|
||||
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 TestRestoreClosesTargetArchiveBeforeReleasingLifecycleLock(t *testing.T) {
|
||||
@@ -179,22 +182,85 @@ func TestRestoreClosesTargetArchiveBeforeReleasingLifecycleLock(t *testing.T) {
|
||||
target = result
|
||||
return result, err
|
||||
}
|
||||
unlockBeforeTargetClose := false
|
||||
deps.acquireLock = func(config.Installation) (restoreLock, error) {
|
||||
return fakeRestoreLock{release: func() {
|
||||
if target.archive != nil && target.archive.file != nil {
|
||||
if _, err := target.archive.file.Stat(); err == nil {
|
||||
unlockBeforeTargetClose = true
|
||||
}
|
||||
}
|
||||
}}, nil
|
||||
}
|
||||
deps.acquireTransaction = lifecycle.AcquireTransaction
|
||||
|
||||
if _, err := restoreWithDependencies(context.Background(), installation, RestoreRequest{Archive: archive, Confirm: true}, deps); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if unlockBeforeTargetClose {
|
||||
t.Fatal("restore released its lifecycle lock before closing the target archive")
|
||||
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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -297,8 +363,8 @@ func TestRestoreLifecycleLockExcludesCompetingTransactionsUntilTerminalCleanup(t
|
||||
caller, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
deps := restoreTestDependencies(t, runner)
|
||||
deps.acquireLock = func(target config.Installation) (restoreLock, error) { return lifecycle.Acquire(target) }
|
||||
deps.checkpointLocked = func(context.Context, config.Installation, CreateRequest) (Result, error) {
|
||||
deps.acquireTransaction = lifecycle.AcquireTransaction
|
||||
deps.checkpoint = func(context.Context, *lifecycle.Transaction, config.Installation, CreateRequest) (Result, error) {
|
||||
gate("checkpoint")
|
||||
return Result{Path: filepath.Join(t.TempDir(), "checkpoint.zip")}, nil
|
||||
}
|
||||
@@ -369,8 +435,8 @@ func competingRestoreAndBackupEntry(installation config.Installation, archive st
|
||||
|
||||
checkpointCalled := false
|
||||
restoreDeps := restoreTestDependencies(t, newBackupRunner(installation, false))
|
||||
restoreDeps.acquireLock = func(target config.Installation) (restoreLock, error) { return lifecycle.Acquire(target) }
|
||||
restoreDeps.checkpointLocked = func(context.Context, config.Installation, CreateRequest) (Result, error) {
|
||||
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
|
||||
}
|
||||
@@ -408,8 +474,8 @@ func TestRestoreCannotApplyAStaleCheckpointOverAnInterleavedRestore(t *testing.T
|
||||
continueCheckpoint := make(chan struct{})
|
||||
firstRunner := newBackupRunner(installation, true)
|
||||
firstDeps := restoreTestDependencies(t, firstRunner)
|
||||
firstDeps.acquireLock = func(target config.Installation) (restoreLock, error) { return lifecycle.Acquire(target) }
|
||||
firstDeps.checkpointLocked = func(context.Context, config.Installation, CreateRequest) (Result, error) {
|
||||
firstDeps.acquireTransaction = lifecycle.AcquireTransaction
|
||||
firstDeps.checkpoint = func(context.Context, *lifecycle.Transaction, config.Installation, CreateRequest) (Result, error) {
|
||||
checkpointState = targetState
|
||||
close(checkpointEntered)
|
||||
<-continueCheckpoint
|
||||
@@ -439,8 +505,8 @@ func TestRestoreCannotApplyAStaleCheckpointOverAnInterleavedRestore(t *testing.T
|
||||
interleavedCheckpoint := false
|
||||
interleavedMutation := false
|
||||
interleavedDeps := restoreTestDependencies(t, newBackupRunner(installation, false))
|
||||
interleavedDeps.acquireLock = func(target config.Installation) (restoreLock, error) { return lifecycle.Acquire(target) }
|
||||
interleavedDeps.checkpointLocked = func(context.Context, config.Installation, CreateRequest) (Result, error) {
|
||||
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
|
||||
}
|
||||
@@ -542,7 +608,7 @@ func TestRestorePreflightFailureDoesNotMutateTarget(t *testing.T) {
|
||||
return PreflightResult{}, preflightErr
|
||||
}
|
||||
checkpointCalls := 0
|
||||
deps.checkpointLocked = func(context.Context, config.Installation, CreateRequest) (Result, error) {
|
||||
deps.checkpoint = func(context.Context, *lifecycle.Transaction, config.Installation, CreateRequest) (Result, error) {
|
||||
checkpointCalls++
|
||||
return Result{}, nil
|
||||
}
|
||||
@@ -567,7 +633,7 @@ func TestRestoreCheckpointFailureDoesNotMutateTarget(t *testing.T) {
|
||||
runner := newBackupRunner(installation, true)
|
||||
deps := restoreTestDependencies(t, runner)
|
||||
checkpointErr := errors.New("checkpoint unavailable")
|
||||
deps.checkpointLocked = func(context.Context, config.Installation, CreateRequest) (Result, error) {
|
||||
deps.checkpoint = func(context.Context, *lifecycle.Transaction, config.Installation, CreateRequest) (Result, error) {
|
||||
return Result{}, checkpointErr
|
||||
}
|
||||
restoredFiles := 0
|
||||
@@ -591,7 +657,7 @@ func TestRestoreFileFailureRollsBackSecretAwareCheckpointBeforeCleanup(t *testin
|
||||
runner := newBackupRunner(installation, true)
|
||||
deps := restoreTestDependencies(t, runner)
|
||||
fileErr := errors.New("cannot restore operator configuration")
|
||||
deps.checkpointLocked = func(context.Context, config.Installation, CreateRequest) (Result, error) {
|
||||
deps.checkpoint = func(context.Context, *lifecycle.Transaction, config.Installation, CreateRequest) (Result, error) {
|
||||
return Result{Path: "/tmp/recovery.zip"}, nil
|
||||
}
|
||||
var events []string
|
||||
@@ -1095,7 +1161,7 @@ func TestRestoreReleasesBarrierOnlyAfterVerifiedRecoveryFromLostResponse(t *test
|
||||
}
|
||||
runner := &commandFailureRunner{fakeBackupRunner: backing, failures: []*commandFailure{&failure}}
|
||||
deps := restoreTestDependencies(t, runner)
|
||||
deps.checkpointLocked = func(context.Context, config.Installation, CreateRequest) (Result, error) {
|
||||
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) {
|
||||
@@ -1322,10 +1388,6 @@ func (runner *lifecycleGateRunner) SessionInventoryScope() string {
|
||||
return runner.fakeBackupRunner.SessionInventoryScope()
|
||||
}
|
||||
|
||||
type fakeRestoreLock struct {
|
||||
release func()
|
||||
}
|
||||
|
||||
type authenticationStateResetRunner struct {
|
||||
args []string
|
||||
result compose.Result
|
||||
@@ -1431,13 +1493,6 @@ func TestVerifyRestoreWorkspaceRejectsInvalidOperatorResults(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func (lock fakeRestoreLock) Release() error {
|
||||
if lock.release != nil {
|
||||
lock.release()
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func equalStrings(got, want []string) bool {
|
||||
if len(got) != len(want) {
|
||||
return false
|
||||
@@ -1456,18 +1511,18 @@ func restoreTestDependencies(t *testing.T, runner archiveRunner) restoreDependen
|
||||
preflight: func(ctx context.Context, installation config.Installation, request PreflightRequest) (PreflightResult, error) {
|
||||
return Preflight(ctx, installation, request, permissivePreflightDependencies())
|
||||
},
|
||||
checkpointLocked: func(context.Context, config.Installation, CreateRequest) (Result, error) {
|
||||
checkpoint: func(context.Context, *lifecycle.Transaction, config.Installation, CreateRequest) (Result, error) {
|
||||
return Result{Path: "/tmp/default-checkpoint.zip"}, nil
|
||||
},
|
||||
prepareRecovery: func(context.Context, config.Installation, string) (PreflightResult, error) {
|
||||
return PreflightResult{}, nil
|
||||
},
|
||||
recover: func(context.Context, config.Installation, PreflightResult, bool) error { return nil },
|
||||
cleanupCheckpoint: func(string) error { return nil },
|
||||
acquireLock: func(config.Installation) (restoreLock, error) { return fakeRestoreLock{}, nil },
|
||||
runner: runner,
|
||||
sleep: func(time.Duration) {},
|
||||
restoreFile: func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error { return nil },
|
||||
recover: func(context.Context, config.Installation, PreflightResult, bool) 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
|
||||
},
|
||||
|
||||
@@ -16,8 +16,10 @@ import (
|
||||
)
|
||||
|
||||
var (
|
||||
ErrLocked = errors.New("another lifecycle operation is already running for this installation")
|
||||
ErrOwnership = errors.New("lifecycle lock ownership changed; refusing to remove it")
|
||||
ErrLocked = errors.New("another lifecycle operation is already running for this installation")
|
||||
ErrOwnership = errors.New("lifecycle lock ownership changed; refusing to remove it")
|
||||
ErrTransactionInactive = errors.New("lifecycle transaction capability is not active")
|
||||
ErrTransactionInstallation = errors.New("lifecycle transaction capability belongs to another installation")
|
||||
)
|
||||
|
||||
const lockFileName = "lifecycle.lock.owner.json"
|
||||
@@ -36,6 +38,14 @@ type Lock struct {
|
||||
released bool
|
||||
}
|
||||
|
||||
// Transaction is an opaque, installation-bound capability for work that must run while a
|
||||
// lifecycle lock remains owned. Its fields are deliberately private so callers can obtain one
|
||||
// only through AcquireTransaction.
|
||||
type Transaction struct {
|
||||
lock *Lock
|
||||
controlDirectory string
|
||||
}
|
||||
|
||||
// Acquire obtains the shared lock used by backup, restore, Pi lifecycle and product updates.
|
||||
func Acquire(installation config.Installation) (*Lock, error) {
|
||||
directory := installation.ControlDirectory()
|
||||
@@ -76,6 +86,39 @@ func Acquire(installation config.Installation) (*Lock, error) {
|
||||
return &Lock{path: path, token: token}, nil
|
||||
}
|
||||
|
||||
// AcquireTransaction obtains a lifecycle lock and returns the capability required by callers
|
||||
// that perform nested work inside the same non-reentrant transaction.
|
||||
func AcquireTransaction(installation config.Installation) (*Transaction, error) {
|
||||
lock, err := Acquire(installation)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &Transaction{
|
||||
lock: lock,
|
||||
controlDirectory: filepath.Clean(installation.ControlDirectory()),
|
||||
}, nil
|
||||
}
|
||||
|
||||
// Verify refuses a nil, released, replaced, or foreign-installation capability before a nested
|
||||
// lifecycle operation can begin.
|
||||
func (transaction *Transaction) Verify(installation config.Installation) error {
|
||||
if transaction == nil || transaction.lock == nil {
|
||||
return ErrTransactionInactive
|
||||
}
|
||||
if transaction.controlDirectory != filepath.Clean(installation.ControlDirectory()) {
|
||||
return ErrTransactionInstallation
|
||||
}
|
||||
return transaction.lock.verifyHeld()
|
||||
}
|
||||
|
||||
// Release relinquishes the lifecycle lock associated with this transaction capability.
|
||||
func (transaction *Transaction) Release() error {
|
||||
if transaction == nil || transaction.lock == nil {
|
||||
return nil
|
||||
}
|
||||
return transaction.lock.Release()
|
||||
}
|
||||
|
||||
// Path returns the installation-private owner-file path for diagnostics and tests.
|
||||
func (lock *Lock) Path() string {
|
||||
if lock == nil {
|
||||
@@ -111,3 +154,26 @@ func (lock *Lock) Release() error {
|
||||
lock.released = true
|
||||
return nil
|
||||
}
|
||||
|
||||
func (lock *Lock) verifyHeld() error {
|
||||
if lock == nil {
|
||||
return ErrTransactionInactive
|
||||
}
|
||||
lock.mu.Lock()
|
||||
defer lock.mu.Unlock()
|
||||
if lock.released {
|
||||
return ErrTransactionInactive
|
||||
}
|
||||
contents, err := os.ReadFile(lock.path)
|
||||
if err != nil {
|
||||
if errors.Is(err, os.ErrNotExist) {
|
||||
return ErrOwnership
|
||||
}
|
||||
return fmt.Errorf("read lifecycle lock owner: %w", err)
|
||||
}
|
||||
var current owner
|
||||
if json.Unmarshal(contents, ¤t) != nil || current.Token == "" || current.Token != lock.token {
|
||||
return ErrOwnership
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -48,3 +48,27 @@ func TestLifecycleLockReleaseDoesNotRemoveAnotherOwnersFile(t *testing.T) {
|
||||
t.Fatalf("foreign lock was removed: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestTransactionCapabilityIsInstallationBoundAndExpiresOnRelease(t *testing.T) {
|
||||
root := t.TempDir()
|
||||
installation := config.Installation{ProjectDirectory: root, Path: filepath.Join(root, "thothii-installation.yaml")}
|
||||
transaction, err := AcquireTransaction(installation)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := transaction.Verify(installation); err != nil {
|
||||
t.Fatalf("Verify() active capability error = %v", err)
|
||||
}
|
||||
|
||||
otherRoot := t.TempDir()
|
||||
other := config.Installation{ProjectDirectory: otherRoot, Path: filepath.Join(otherRoot, "thothii-installation.yaml")}
|
||||
if err := transaction.Verify(other); !errors.Is(err, ErrTransactionInstallation) {
|
||||
t.Fatalf("Verify() for another installation error = %v, want ErrTransactionInstallation", err)
|
||||
}
|
||||
if err := transaction.Release(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := transaction.Verify(installation); !errors.Is(err, ErrTransactionInactive) {
|
||||
t.Fatalf("Verify() after Release() error = %v, want ErrTransactionInactive", err)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user