Files
ThothII/tools/tht/internal/backup/restore.go
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

452 lines
18 KiB
Go

package backup
import (
"archive/zip"
"context"
"errors"
"fmt"
"io"
"time"
"github.com/aritmolab/thothii/tools/tht/internal/authconfig"
"github.com/aritmolab/thothii/tools/tht/internal/config"
"github.com/aritmolab/thothii/tools/tht/internal/lifecycle"
)
var ErrRestoreConfirmationRequired = errors.New("restore requires --yes")
// RestoreRequest describes the deliberately-confirmed archive restoration.
type RestoreRequest struct {
Archive string
Confirm bool
Drain bool
}
// RestoreResult records the final service state. Secret-bearing recovery checkpoints are always
// destroyed internally and are therefore never exposed to callers.
type RestoreResult struct {
// Deprecated: always empty. Recovery checkpoints are private transaction internals.
Checkpoint string
Restarted bool
Verified bool
}
type restoreVerify func(context.Context, config.Installation, archiveRunner) error
type authProjectionRestoreTransaction interface {
PublishCanonical() (authconfig.ProjectionStatus, error)
RestorePriorIfCanonicalUnchanged() error
Close() error
}
type restoreDependencies struct {
preflight func(context.Context, config.Installation, PreflightRequest) (PreflightResult, 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, *stagedArchive, bool, authProjectionRestoreTransaction) error
requireAuthProjection func() error
beginAuthProjection func(context.Context, config.Installation) (authProjectionRestoreTransaction, error)
cleanupCheckpoint func(string) 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
restoreVolume func(context.Context, config.Installation, VolumeMetadata, io.Reader) error
resetAuthenticationState func(context.Context, config.Installation, archiveRunner) error
verify map[string]restoreVerify
}
// restoreTransactionState is the restore admission state machine. The durable maintenance
// barrier is released only after the target archive has been verified, or after a separately
// verified checkpoint recovery. Every mutable Docker command is tracked before it is invoked.
type restoreTransactionState struct {
wasRunning bool
maintenanceAttempted bool
stopAttempted bool
mutated bool
verified bool
recovered bool
}
func (state restoreTransactionState) recoveryRequired(resultErr error) bool {
return resultErr != nil && state.mutated && !state.verified
}
func (state restoreTransactionState) mayDeactivateMaintenance() bool {
return !state.mutated || state.verified || state.recovered
}
// Restore runs the host transaction through the same concrete Docker/filesystem boundaries used
// by backup creation. The injectable core below exists only to make every failure boundary
// deterministic in tests.
func Restore(ctx context.Context, installation config.Installation, request RestoreRequest) (RestoreResult, error) {
return restoreWithDependencies(ctx, installation, request, productionRestoreDependencies(installation))
}
func restoreWithDependencies(ctx context.Context, installation config.Installation, request RestoreRequest, deps restoreDependencies) (result RestoreResult, resultErr error) {
if !request.Confirm {
return RestoreResult{}, ErrRestoreConfirmationRequired
}
if request.Archive == "" {
return RestoreResult{}, errors.New("restore archive is required")
}
if deps.preflight == nil || deps.checkpoint == nil || deps.prepareRecovery == nil || deps.recover == nil || deps.requireAuthProjection == 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")
}
// 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.
transaction, err := deps.acquireTransaction(installation)
if err != nil {
return result, err
}
defer func() {
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()
authRestoreRequired := installation.HasRuntimeAuthProjection() && manifestArchivesAuthentication(preflight.Manifest)
if authRestoreRequired {
if err := deps.requireAuthProjection(); err != nil {
return result, err
}
}
checkpoint, err := deps.checkpoint(ctx, transaction, installation, CreateRequest{IncludeSecrets: true, Confirm: true})
if err != nil {
return result, fmt.Errorf("create recovery checkpoint: %w", err)
}
recovery, err := deps.prepareRecovery(ctx, installation, checkpoint.Path)
if err != nil {
cleanupErr := deps.cleanupCheckpoint(checkpoint.Path)
return result, errors.Join(fmt.Errorf("validate recovery checkpoint: %w", err), cleanupErr)
}
defer recovery.CloseArchive()
// Both immutable archives are staged while the lifecycle transaction is held and before
// maintenance, service stops, or destination writes. Keeping both files open reserves their
// combined staging capacity, so checkpoint recovery never needs a new allocation after a
// candidate mutation has begun.
candidateStage, err := preflight.StageArchive(ctx)
if err != nil {
cleanupErr := deps.cleanupCheckpoint(checkpoint.Path)
return result, errors.Join(err, cleanupErr)
}
defer func() {
if closeErr := candidateStage.Close(); closeErr != nil {
result = RestoreResult{}
resultErr = errors.Join(resultErr, closeErr)
}
}()
recoveryStage, err := recovery.stageArchiveAlongside(ctx, candidateStage)
if err != nil {
cleanupErr := deps.cleanupCheckpoint(checkpoint.Path)
return result, errors.Join(err, cleanupErr)
}
defer func() {
if closeErr := recoveryStage.Close(); closeErr != nil {
result = RestoreResult{}
resultErr = errors.Join(resultErr, closeErr)
}
}()
if err := ensureCombinedRestoreCapacity(preflight, recovery); err != nil {
cleanupErr := deps.cleanupCheckpoint(checkpoint.Path)
return result, errors.Join(err, cleanupErr)
}
var authTransaction authProjectionRestoreTransaction
if authRestoreRequired {
if deps.beginAuthProjection == nil {
cleanupErr := deps.cleanupCheckpoint(checkpoint.Path)
return result, errors.Join(errors.New("restore authentication projection dependency is unavailable"), cleanupErr)
}
authTransaction, err = deps.beginAuthProjection(ctx, installation)
if err != nil {
cleanupErr := deps.cleanupCheckpoint(checkpoint.Path)
return result, errors.Join(errors.New("restore authentication projection could not be blocked"), cleanupErr)
}
}
state := restoreTransactionState{}
defer func() {
if state.recoveryRequired(resultErr) {
recoveryContext, cancel := boundedCleanupContext()
recoveryErr := deps.recover(recoveryContext, installation, recovery, recoveryStage, state.wasRunning, authTransaction)
cancel()
if recoveryErr != nil {
resultErr = errors.Join(resultErr, fmt.Errorf("restore recovery checkpoint: %w", recoveryErr))
if recoveryReachedVerifiedState(recoveryErr) {
state.recovered = true
state.stopAttempted = false
}
} else {
state.recovered = true
state.stopAttempted = false
}
}
checkpointCleanupSucceeded := true
var cleanupErr error
if checkpointErr := deps.cleanupCheckpoint(checkpoint.Path); checkpointErr != nil {
checkpointCleanupSucceeded = false
cleanupErr = errors.Join(cleanupErr, fmt.Errorf("destroy recovery checkpoint: %w", checkpointErr))
}
authCleanupSucceeded := true
if authTransaction != nil {
if resultErr != nil && !state.mutated {
if restoreErr := authTransaction.RestorePriorIfCanonicalUnchanged(); restoreErr != nil {
authCleanupSucceeded = false
cleanupErr = errors.Join(cleanupErr, errors.New("restore authentication projection could not restore its prior selector"))
}
}
if closeErr := authTransaction.Close(); closeErr != nil {
authCleanupSucceeded = false
cleanupErr = errors.Join(cleanupErr, errors.New("restore authentication projection transaction could not be closed"))
}
authTransaction = nil
}
if !state.maintenanceAttempted {
if cleanupErr != nil {
result = RestoreResult{}
resultErr = errors.Join(resultErr, cleanupErr)
}
return
}
// A recovery checkpoint remains secret-bearing transaction state. Do not reopen
// admissions until it has been safely deleted, even if the restored target verified.
if !checkpointCleanupSucceeded {
result = RestoreResult{}
resultErr = errors.Join(resultErr, cleanupErr)
return
}
restartCompleted := true
if state.wasRunning && state.stopAttempted && state.mayDeactivateMaintenance() {
startErr, started := retryBoundedCleanup(func(cleanupContext context.Context) error {
return composeStartAndVerify(cleanupContext, installation, deps.runner)
})
if startErr != nil {
cleanupErr = errors.Join(cleanupErr, fmt.Errorf("restore maintenance cleanup restart: %w", startErr))
}
if started {
state.stopAttempted = false
} else {
restartCompleted = false
}
}
// A failed checkpoint recovery deliberately leaves admissions blocked. Starting or
// deactivating at that point would expose an unverified, possibly partial restore.
if state.mayDeactivateMaintenance() && restartCompleted && authCleanupSucceeded {
deactivateErr, deactivated := retryBoundedCleanup(func(cleanupContext context.Context) error {
return maintenance(cleanupContext, installation, deps.runner, false)
})
if deactivateErr != nil {
cleanupErr = errors.Join(cleanupErr, fmt.Errorf("restore maintenance cleanup: %w", deactivateErr))
}
if deactivated {
state.maintenanceAttempted = false
}
}
if cleanupErr != nil {
result = RestoreResult{}
resultErr = errors.Join(resultErr, cleanupErr)
}
}()
wasRunning, err := installationRunning(ctx, installation, deps.runner)
if err != nil {
return result, err
}
state.wasRunning = wasRunning
if state.wasRunning {
// The activation command may take effect even when its response is lost. Track the attempt,
// not merely a successful return, so every subsequent path compensates from durable state.
state.maintenanceAttempted = true
if err := maintenance(ctx, installation, deps.runner, true); err != nil {
return result, err
}
if err := waitForNoActiveSessions(ctx, installation, deps.runner, request.Drain, deps.sleep); err != nil {
return result, err
}
// Compose may stop the core and then lose its response. Cleanup must therefore restart after
// any stop attempt, including a command that returns an error.
state.stopAttempted = true
if err := runCompose(ctx, installation, deps.runner, "stop"); err != nil {
return result, err
}
}
state.mutated = true
if err := restoreVerifiedEntries(ctx, installation, preflight, candidateStage.file, deps.restoreFile, deps.restoreVolume); err != nil {
return result, err
}
if err := deps.resetAuthenticationState(ctx, installation, deps.runner); err != nil {
return result, fmt.Errorf("reset authentication state: %w", err)
}
if authTransaction != nil {
if err := publishRestoredAuthentication(authTransaction); err != nil {
return result, err
}
}
if state.wasRunning {
if err := composeStartAndVerify(ctx, installation, deps.runner); err != nil {
return result, err
}
state.stopAttempted = false
result.Restarted = true
}
if err := verifyRestoreTransaction(ctx, installation, deps); err != nil {
return result, err
}
state.verified = true
result.Verified = true
return result, nil
}
func manifestArchivesAuthentication(manifest Manifest) bool {
for _, entry := range manifest.Entries {
if entry.Archived && entry.Kind == EntryExternalSecret && entry.Owner == "authentication-configuration" {
return true
}
}
return false
}
func publishRestoredAuthentication(transaction authProjectionRestoreTransaction) error {
status, err := transaction.PublishCanonical()
if err != nil || status.State != "ready" || !status.Equal {
if err != nil {
return fmt.Errorf("publish restored authentication projection: %w", err)
}
return errors.New("publish restored authentication projection")
}
return nil
}
func ensureCombinedRestoreCapacity(candidate, recovery PreflightResult) error {
if candidate.freeBytes == nil || candidate.stagingRoot == "" || candidate.stagingRoot != recovery.stagingRoot {
return errors.New("candidate and recovery archives do not share controlled restore staging")
}
if recovery.RequiredBytes > ^uint64(0)-candidate.RequiredBytes {
return errors.New("combined candidate and recovery restore capacity exceeds supported size")
}
required := candidate.RequiredBytes + recovery.RequiredBytes
freeBytes, err := candidate.freeBytes(candidate.stagingRoot)
if err != nil {
return errors.New("check combined candidate and recovery staging capacity")
}
if freeBytes < required {
return fmt.Errorf("insufficient free disk space for combined candidate and recovery staging: need %d bytes, have %d", required, freeBytes)
}
return nil
}
func verifyRestoreTransaction(ctx context.Context, installation config.Installation, deps restoreDependencies) error {
for _, name := range []string{"health", "doctor", "pi", "workspace"} {
check := deps.verify[name]
if check == nil {
return fmt.Errorf("restore verification %q is unavailable", name)
}
if err := check(ctx, installation, deps.runner); err != nil {
return fmt.Errorf("restore verification %s: %w", name, err)
}
}
return nil
}
type verifiedRecoveryError struct{ err error }
func (err *verifiedRecoveryError) Error() string { return err.err.Error() }
func (err *verifiedRecoveryError) Unwrap() error { return err.err }
func recoveryReachedVerifiedState(err error) bool {
var verifiedErr *verifiedRecoveryError
return errors.As(err, &verifiedErr)
}
func restoreVerifiedEntries(
ctx context.Context,
installation config.Installation,
preflight PreflightResult,
archive io.ReaderAt,
restoreFile func(context.Context, config.Installation, ArchiveEntryMetadata, io.Reader) error,
restoreVolume func(context.Context, config.Installation, VolumeMetadata, io.Reader) error,
) error {
reader, err := zip.NewReader(archive, preflight.ArchiveSize)
if err != nil {
return fmt.Errorf("read verified restore archive: %w", err)
}
members := make(map[string]*zip.File, len(reader.File))
for _, member := range reader.File {
members[member.Name] = member
}
for _, entry := range preflight.Entries {
member := members[entry.Path]
if member == nil {
return fmt.Errorf("verified archive is missing %q", entry.Path)
}
stream, openErr := member.Open()
if openErr != nil {
return fmt.Errorf("open verified archive member %q: %w", entry.Path, openErr)
}
var restoreErr error
if entry.Kind == EntryVolume {
volume, found := restoreVolumeMetadata(preflight.Manifest, entry.LogicalName)
if !found {
_ = stream.Close()
return errors.New("verified volume metadata is incomplete")
}
restoreErr = restoreVolume(ctx, installation, volume, stream)
} else {
restoreErr = restoreFile(ctx, installation, entry, stream)
}
closeErr := stream.Close()
if restoreErr != nil {
return restoreErr
}
if closeErr != nil {
return closeErr
}
}
return nil
}
func restoreVolumeMetadata(manifest Manifest, logicalName string) (VolumeMetadata, bool) {
if logicalName == "" {
return VolumeMetadata{}, false
}
for _, volume := range manifest.Volumes {
if volume.LogicalName == logicalName {
return volume, true
}
}
return VolumeMetadata{}, false
}
// resetAuthenticationState clears browser sessions and pending OIDC transactions without touching
// installation-global auth.yaml or users.yaml. A root-scoped one-shot repairs ownership and mode
// before clearing children; links and malformed roots are rejected before any recursive removal.
func resetAuthenticationState(ctx context.Context, installation config.Installation, runner archiveRunner) error {
result, err := runner.Run(ctx, installation.ComposeArgs(
"run", "--rm", "--no-deps", "--no-TTY", "--user", "0:0", "--entrypoint", "sh", "core", "-ceu",
"test ! -L /data/auth && { test ! -e /data/auth || test -d /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 && chmod 0700 /data/auth /data/auth/sessions /data/auth/oidc && chown 10001:10001 /data/auth /data/auth/sessions /data/auth/oidc && test -z \"$(find /data/auth/sessions /data/auth/oidc -mindepth 1 -print -quit)\"",
), nil)
if err != nil {
return dockerError("reset authentication state", result, err)
}
if result.ExitCode != 0 {
return dockerError("reset authentication state", result, errors.New("Compose returned a nonzero exit status"))
}
return nil
}