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

477 lines
19 KiB
Go

package backup
import (
"context"
"crypto/rand"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"errors"
"fmt"
"io"
"os"
"path/filepath"
"regexp"
"strconv"
"strings"
"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/doctor"
"github.com/aritmolab/thothii/tools/tht/internal/lifecycle"
"github.com/aritmolab/thothii/tools/tht/internal/pi"
"github.com/aritmolab/thothii/tools/tht/internal/safeio"
"github.com/aritmolab/thothii/tools/tht/internal/service"
)
func productionRestoreDependencies(installation config.Installation) restoreDependencies {
runner := hostRunner{runner: compose.NewRunner(""), binary: "docker", profile: installation.Profile}
deps := restoreDependencies{
preflight: func(ctx context.Context, target config.Installation, request PreflightRequest) (PreflightResult, error) {
return Preflight(ctx, target, request, PreflightDependencies{
FreeBytes: restoreFreeBytes,
CheckOwnershipPermissions: validateRestoreTargets,
CheckVolumeMapping: func(ctx context.Context, target config.Installation, manifest Manifest) error {
return validateRestoreVolumes(ctx, target, manifest, runner)
},
CheckImageConfigCompatibility: func(ctx context.Context, target config.Installation, manifest Manifest) error {
return validateRestoreImages(ctx, target, manifest, runner)
},
})
},
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 createWithDependenciesTransaction(ctx, transaction, target, request, productionCreateDependencies(target))
},
cleanupCheckpoint: cleanupRecoveryCheckpoint,
acquireTransaction: lifecycle.AcquireTransaction,
requireAuthProjection: requireAuthProjectionRestorePrivilege,
beginAuthProjection: func(ctx context.Context, target config.Installation) (authProjectionRestoreTransaction, error) {
projection := target.RuntimeAuthProjection()
if projection == nil {
return nil, errors.New("runtime authentication projection is unavailable")
}
return authconfig.BeginExternalProjectionTransaction(ctx, target.AuthenticationDirectory(), authconfig.ProjectionSpec{
RuntimeRoot: projection.Directory,
UID: projection.UID,
GID: projection.GID,
})
},
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 {
if err == nil {
err = errors.New("Docker volume helper returned a nonzero exit status")
}
return fmt.Errorf("restore volume %s: %w", volume.LogicalName, dockerError("stream volume", result, err))
}
return nil
},
resetAuthenticationState: resetAuthenticationState,
verify: map[string]restoreVerify{
"health": verifyRestoreHealth,
"doctor": verifyRestoreDoctor,
"pi": verifyRestorePi,
"workspace": verifyRestoreWorkspace,
},
}
deps.prepareRecovery = func(ctx context.Context, target config.Installation, path string) (PreflightResult, error) {
return deps.preflight(ctx, target, PreflightRequest{Archive: path, Confirm: true, AllowExternalSecrets: true})
}
deps.recover = func(ctx context.Context, target config.Installation, recovery PreflightResult, staged *stagedArchive, wasRunning bool, transaction authProjectionRestoreTransaction) error {
return recoverRestoreTransaction(ctx, target, recovery, staged, wasRunning, deps, transaction)
}
return deps
}
func cleanupRecoveryCheckpoint(path string) error {
if err := safeio.RemoveCanonicalPrivateRegular(path); err != nil {
return errors.New("private recovery checkpoint could not be destroyed safely")
}
return nil
}
func recoverRestoreTransaction(ctx context.Context, installation config.Installation, recovery PreflightResult, staged *stagedArchive, wasRunning bool, deps restoreDependencies, transaction authProjectionRestoreTransaction) (resultErr error) {
if staged == nil || staged.file == nil {
return errors.New("recovery checkpoint was not staged before restore mutation")
}
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
// idempotent command so no recovery mutation starts while the candidate core is running.
if retryErr := runCompose(ctx, installation, deps.runner, "stop"); retryErr != nil {
return errors.Join(resultErr, retryErr)
}
}
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
// authentication state is deliberately reset again so neither sessions nor OIDC transactions
// survive a failed restore attempt.
if err := deps.resetAuthenticationState(ctx, installation, deps.runner); err != nil {
return errors.Join(resultErr, err)
}
if transaction != nil {
if err := publishRestoredAuthentication(transaction); err != nil {
return errors.Join(resultErr, err)
}
}
if wasRunning {
if err := composeStartAndVerify(ctx, installation, deps.runner); err != nil {
resultErr = errors.Join(resultErr, err)
// As with stop, a start can succeed while the client loses its response. A successful
// retry includes its own health check before recovery verification proceeds.
if retryErr := composeStartAndVerify(ctx, installation, deps.runner); retryErr != nil {
return errors.Join(resultErr, retryErr)
}
}
}
if err := verifyRestoreTransaction(ctx, installation, deps); err != nil {
return errors.Join(resultErr, err)
}
if resultErr != nil {
// A prior response was lost, but the checkpoint has been restored and fully verified. The
// caller may safely remove maintenance while still returning every observed error.
return &verifiedRecoveryError{err: resultErr}
}
return nil
}
func restoreCheckpointPath(installation config.Installation, now time.Time) (string, error) {
suffix := make([]byte, 8)
if _, err := rand.Read(suffix); err != nil {
return "", errors.New("recovery checkpoint name is unavailable")
}
name := fmt.Sprintf("restore-checkpoint-%s-%s.zip", now.Format("20060102T150405.000000000Z"), hex.EncodeToString(suffix))
return filepath.Join(installation.ControlDirectory(), name), nil
}
func validateRestoreTargets(_ context.Context, installation config.Installation, manifest Manifest) error {
for _, entry := range manifest.Entries {
if entry.Kind == EntryVolume {
continue
}
if !entry.Archived {
if entry.Kind == EntrySecretReference || entry.Kind == EntryPreservationReference {
if err := validateRestoreReference(installation, entry); err != nil {
return err
}
}
continue
}
metadata := ArchiveEntryMetadata{
Path: entry.Path, Kind: entry.Kind, Owner: entry.Owner, LogicalName: entry.LogicalName,
SourcePath: entry.SourcePath, Size: entry.Size, SHA256: entry.SHA256, Mode: entry.Mode,
Sensitive: entry.Sensitive,
}
target, err := restoreFileTarget(installation, metadata)
if err != nil || !safeRestoreParent(target) {
return errors.New("restore target ownership or permissions are invalid")
}
}
return nil
}
func validateRestoreReference(installation config.Installation, entry Entry) error {
metadata := ArchiveEntryMetadata{
Path: entry.Path, Kind: entry.Kind, Owner: entry.Owner, SourcePath: entry.SourcePath,
Size: entry.Size, SHA256: entry.SHA256, Mode: entry.Mode, Sensitive: entry.Sensitive,
}
if entry.Kind == EntryPreservationReference {
return nil
}
if _, err := restoreExternalTarget(installation, metadata); err != nil {
return errors.New("restore external prerequisite is invalid")
}
contents, err := safeio.ReadCanonicalRegular(entry.SourcePath, entry.Size)
if err != nil || int64(len(contents)) != entry.Size {
return errors.New("restore external prerequisite is unavailable or unsafe")
}
digest := sha256.Sum256(contents)
if "sha256:"+hex.EncodeToString(digest[:]) != entry.SHA256 {
return errors.New("restore external prerequisite has changed")
}
return nil
}
func validateRestoreVolumes(ctx context.Context, installation config.Installation, manifest Manifest, runner archiveRunner) error {
if len(manifest.Volumes) != len(requiredBackupVolumes(installation)) {
return errors.New("backup volume set is incomplete")
}
rendered, err := renderedConfiguration(ctx, installation, runner)
if err != nil {
return err
}
current, err := inspectRequiredVolumes(ctx, installation, runner, rendered)
if err != nil {
return err
}
archived := make(map[string]VolumeMetadata, len(manifest.Volumes))
for _, volume := range manifest.Volumes {
archived[volume.LogicalName] = volume
}
for _, volume := range current {
previous, found := archived[volume.LogicalName]
if !found || previous.Name != volume.Name || previous.Driver != volume.Driver {
return errors.New("backup volume ownership does not match the installation")
}
if volume.Labels["com.docker.compose.project"] != installation.ProjectName() {
return errors.New("current volume is not owned by the installation")
}
}
return nil
}
func validateRestoreImages(ctx context.Context, installation config.Installation, manifest Manifest, runner archiveRunner) error {
rendered, err := renderedConfiguration(ctx, installation, runner)
if err != nil {
return err
}
for _, image := range manifest.Images {
serviceDefinition, found := rendered.Services[image.Service]
if !found || serviceDefinition.Image == "" || serviceDefinition.Image != image.Reference {
return errors.New("backup image configuration does not match the installation")
}
}
return nil
}
func restoreFilePayload(_ context.Context, installation config.Installation, entry ArchiveEntryMetadata, input io.Reader) error {
if entry.Size < 0 || uint64(entry.Size) > defaultPreflightMaxUncompressedBytes {
return errors.New("restore file size is invalid")
}
contents, err := io.ReadAll(io.LimitReader(input, entry.Size+1))
if err != nil || int64(len(contents)) != entry.Size {
return errors.New("restore file payload is invalid")
}
target, err := restoreFileTarget(installation, entry)
if err != nil {
return err
}
if err := replaceRestoreFile(target, contents, os.FileMode(entry.Mode)); err != nil {
return errors.New("restore file could not be replaced safely")
}
return nil
}
func restoreFileTarget(installation config.Installation, entry ArchiveEntryMetadata) (string, error) {
if entry.Kind == EntryExternalSecret {
return restoreExternalTarget(installation, entry)
}
if entry.Kind != EntryFile {
return "", errors.New("restore file kind is unsupported")
}
switch entry.Path {
case "configuration/installation/thothii-installation.yaml":
return installation.Path, nil
case "configuration/environment/operator.env":
return installation.EnvFile, nil
case "configuration/pi/models.json":
return filepath.Join(installation.ProjectDirectory, "deploy", "pi", "models.json"), nil
case "configuration/pi/settings.json":
return filepath.Join(installation.ProjectDirectory, "deploy", "pi", "settings.json"), nil
case "configuration/generated/current-image.yaml":
return installation.CurrentImageOverridePath(), nil
}
if strings.HasPrefix(entry.Path, "configuration/overrides/") {
name := strings.TrimPrefix(entry.Path, "configuration/overrides/")
indexText, base, found := strings.Cut(name, "-")
index, parseErr := strconv.Atoi(indexText)
if !found || parseErr != nil || len(indexText) != 2 || index < 0 || index >= len(installation.Overrides) || filepath.Base(installation.Overrides[index]) != base {
return "", errors.New("restore override target is invalid")
}
return installation.Overrides[index], nil
}
if strings.HasPrefix(entry.Owner, "preservation-root:") && strings.HasPrefix(entry.Path, "preservation/") {
variable := strings.TrimPrefix(entry.Owner, "preservation-root:")
allowed := variable == "THT_DATA_ROOT" || variable == "THT_PI_STATE_ROOT" || variable == "THT_WORKSPACE_REGISTRY_ROOT"
parts := strings.SplitN(entry.Path, "/", 3)
root, rootErr := installation.EnvironmentValue(variable)
if !allowed || len(parts) != 3 || rootErr != nil || root == "" {
return "", errors.New("restore preservation target is invalid")
}
target := filepath.Join(root, filepath.FromSlash(parts[2]))
relative, relErr := filepath.Rel(root, target)
if relErr != nil || relative == ".." || strings.HasPrefix(relative, ".."+string(filepath.Separator)) {
return "", errors.New("restore preservation target escapes its root")
}
return target, nil
}
return "", errors.New("restore file target is not declared")
}
func restoreExternalTarget(installation config.Installation, entry ArchiveEntryMetadata) (string, error) {
if entry.SourcePath == "" || filepath.Clean(entry.SourcePath) != entry.SourcePath || !filepath.IsAbs(entry.SourcePath) {
return "", errors.New("restore external target is invalid")
}
if entry.Owner == "external-secret" {
paths, err := installation.SecretFiles()
if err != nil {
return "", errors.New("restore external secret declarations are unavailable")
}
for _, path := range paths {
if path == entry.SourcePath {
return path, nil
}
}
return "", errors.New("restore external secret target is not declared")
}
if entry.Owner == "authentication-configuration" {
for _, name := range []string{"auth.yaml", "users.yaml"} {
path := filepath.Join(installation.AuthenticationDirectory(), name)
if entry.SourcePath == path {
return path, nil
}
}
}
return "", errors.New("restore external target owner is invalid")
}
func safeRestoreParent(target string) bool {
if target == "" || !filepath.IsAbs(target) || filepath.Clean(target) != target {
return false
}
parent := filepath.Dir(target)
resolved, err := filepath.EvalSymlinks(parent)
if err != nil || resolved != parent {
return false
}
info, err := os.Stat(parent)
if err != nil || !info.IsDir() {
return false
}
if targetInfo, err := os.Lstat(target); err == nil {
return targetInfo.Mode().IsRegular() && targetInfo.Mode()&os.ModeSymlink == 0
} else {
return errors.Is(err, os.ErrNotExist)
}
}
func volumeRestoreCommand(volume string) []string {
return []string{
"run", "--rm", "--interactive", "--network", "none", "--mount", "type=volume,src=" + volume + ",dst=/target",
helperImage, "sh", "-ceu",
"rm -rf -- /target/* /target/.[!.]* /target/..?*; tar --numeric-owner -C /target -xf -",
}
}
func restoreVerificationRunning(ctx context.Context, installation config.Installation, runner archiveRunner) (bool, error) {
return installationRunning(ctx, installation, runner)
}
func verifyRestoreHealth(ctx context.Context, installation config.Installation, runner archiveRunner) error {
running, err := restoreVerificationRunning(ctx, installation, runner)
if err != nil || !running {
return err
}
return service.WaitForHealthy(ctx, installation, runner)
}
func verifyRestoreDoctor(ctx context.Context, installation config.Installation, runner archiveRunner) error {
running, err := restoreVerificationRunning(ctx, installation, runner)
if err != nil {
return err
}
if !running {
result, configErr := runner.Run(ctx, installation.ComposeArgs("config", "--quiet"), nil)
if configErr != nil {
return dockerError("verify restored Compose configuration", result, configErr)
}
diagnostics, diagnosticErr := authconfig.Check(ctx, installation, runner, false, false)
if diagnosticErr != nil || !diagnostics.Ready {
if diagnosticErr != nil {
return diagnosticErr
}
return errors.New("authentication diagnostics did not pass after restore")
}
return nil
}
report, err := doctor.Run(ctx, installation, runner)
if err != nil {
return err
}
if !report.OK {
failed := make([]string, 0, len(report.Checks))
for _, check := range report.Checks {
if check.Status == doctor.StatusFailed {
failed = append(failed, check.Name)
}
}
if len(failed) == 0 {
return errors.New("aggregate doctor did not pass after restore")
}
return fmt.Errorf("aggregate doctor failed checks: %s", strings.Join(failed, ","))
}
return nil
}
func verifyRestorePi(ctx context.Context, installation config.Installation, runner archiveRunner) error {
running, err := restoreVerificationRunning(ctx, installation, runner)
if err != nil || !running {
if err != nil {
return err
}
result, runErr := runner.Run(ctx, installation.ComposeArgs(
"run", "--rm", "--no-deps", "--no-TTY", "core", "pi", "--version",
), nil)
if runErr != nil || strings.TrimSpace(result.Stdout) == "" {
if runErr == nil {
runErr = errors.New("Pi version probe returned no version")
}
return dockerError("verify restored Pi runtime", result, runErr)
}
return nil
}
return pi.Doctor(ctx, compose.InstallationRunner{Installation: installation, Runner: runner})
}
func verifyRestoreWorkspace(ctx context.Context, installation config.Installation, runner archiveRunner) error {
running, err := restoreVerificationRunning(ctx, installation, runner)
if err != nil {
return err
}
command := []string{"exec", "-T", "core", "node", "/app/backend/dist/operator-command.js", "workspace-integrity"}
if !running {
command = []string{"run", "--rm", "--no-deps", "--no-TTY", "core", "node", "/app/backend/dist/operator-command.js", "workspace-integrity"}
}
result, err := runner.Run(ctx, installation.ComposeArgs(command...), nil)
if err != nil || result.ExitCode != 0 {
return dockerError("validate restored workspace registry", result, err)
}
if len(result.Stdout) == 0 || len(result.Stdout) > 4096 {
return errors.New("restored workspace registry returned an invalid result")
}
var payload struct {
Ready bool `json:"ready"`
State string `json:"state"`
Workspaces int `json:"workspaces"`
Fingerprint string `json:"fingerprint"`
}
decoder := json.NewDecoder(strings.NewReader(result.Stdout))
decoder.DisallowUnknownFields()
if decodeErr := decoder.Decode(&payload); decodeErr != nil {
return errors.New("restored workspace registry returned an invalid result")
}
var trailing any
if decodeErr := decoder.Decode(&trailing); !errors.Is(decodeErr, io.EOF) {
return errors.New("restored workspace registry returned an invalid result")
}
if !payload.Ready || payload.Workspaces < 0 ||
!regexp.MustCompile(`^sha256:[0-9a-f]{64}$`).MatchString(payload.Fingerprint) ||
(payload.State != "active" && payload.State != "uninitialized") ||
(payload.State == "uninitialized" && payload.Workspaces != 0) {
return errors.New("restored workspace registry did not pass integrity validation")
}
return nil
}