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) }, }) }, checkpointLocked: func(ctx context.Context, 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)) }, cleanupCheckpoint: cleanupRecoveryCheckpoint, acquireLock: func(target config.Installation) (restoreLock, error) { return lifecycle.Acquire(target) }, 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, wasRunning bool) error { return recoverRestoreTransaction(ctx, target, recovery, wasRunning, deps) } 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, wasRunning bool, deps restoreDependencies) error { var 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 // 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) } } archive, err := recovery.RevalidateArchive() if err != nil { return errors.Join(resultErr, err) } if err := restoreVerifiedEntries(ctx, installation, recovery, archive, 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 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 }