Implement restore planning for the restore subcommand
This commit is contained in:
344
internal/app/restore_plan.go
Normal file
344
internal/app/restore_plan.go
Normal file
@@ -0,0 +1,344 @@
|
||||
package app
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
"path"
|
||||
"path/filepath"
|
||||
"sort"
|
||||
"strings"
|
||||
|
||||
"gitea.maximumdirect.net/eric/narratio/internal/adapters/storage"
|
||||
"gitea.maximumdirect.net/eric/narratio/internal/artifacts"
|
||||
"gitea.maximumdirect.net/eric/narratio/internal/config"
|
||||
)
|
||||
|
||||
// RestoreActionKind identifies one restore planner action.
|
||||
type RestoreActionKind string
|
||||
|
||||
const (
|
||||
RestoreActionDownload RestoreActionKind = "download"
|
||||
RestoreActionSkipSame RestoreActionKind = "skip_same"
|
||||
RestoreActionConflict RestoreActionKind = "conflict"
|
||||
)
|
||||
|
||||
// RestoreAction is one deterministic planner action.
|
||||
type RestoreAction struct {
|
||||
Kind RestoreActionKind
|
||||
RemoteKey string
|
||||
LocalRelativePath string
|
||||
LocalPath string
|
||||
Size int64
|
||||
ETag string
|
||||
ExistsLocal bool
|
||||
SameLocal bool
|
||||
Conflict bool
|
||||
Reason string
|
||||
}
|
||||
|
||||
// RestorePlan is the deterministic output of restore planning.
|
||||
type RestorePlan struct {
|
||||
Actions []RestoreAction
|
||||
DownloadCount int
|
||||
SkipSameCount int
|
||||
ConflictCount int
|
||||
}
|
||||
|
||||
// RestorePlanOptions control restore planning scope and classification.
|
||||
type RestorePlanOptions struct {
|
||||
IncludeAudio bool
|
||||
Force bool
|
||||
DryRun bool
|
||||
}
|
||||
|
||||
func buildRestorePlan(ctx context.Context, cfg *config.Config, current *RemoteCurrentState, store storage.ObjectStore, opts RestorePlanOptions) (*RestorePlan, error) {
|
||||
if cfg == nil || cfg.Pipeline == nil || cfg.Session == nil {
|
||||
return nil, fmt.Errorf("resolved config with pipeline/session is required")
|
||||
}
|
||||
if current == nil {
|
||||
return nil, fmt.Errorf("remote current state is required")
|
||||
}
|
||||
if store == nil {
|
||||
return nil, fmt.Errorf("remote object store is required")
|
||||
}
|
||||
prefix := normalizeRemoteKey(current.SessionPrefix)
|
||||
if strings.TrimSpace(prefix) == "" {
|
||||
return nil, fmt.Errorf("remote session prefix is required")
|
||||
}
|
||||
if !strings.HasSuffix(prefix, "/") {
|
||||
prefix += "/"
|
||||
}
|
||||
|
||||
sessionPaths := artifacts.NewLocalStore(cfg.Pipeline.Workspace.Root).SessionPathsFor(cfg.Session.Campaign, cfg.Session.SessionID)
|
||||
objects, err := store.List(ctx, prefix)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("list remote session objects under %q: %w", prefix, err)
|
||||
}
|
||||
|
||||
candidates := make(map[string]storage.ObjectInfo, len(objects)+1)
|
||||
for _, obj := range objects {
|
||||
key := normalizeRemoteKey(obj.Key)
|
||||
if key == "" {
|
||||
continue
|
||||
}
|
||||
obj.Key = key
|
||||
candidates[key] = obj
|
||||
}
|
||||
if strings.TrimSpace(current.CurrentManifestKey) != "" {
|
||||
key := normalizeRemoteKey(current.CurrentManifestKey)
|
||||
if _, ok := candidates[key]; !ok {
|
||||
candidates[key] = storage.ObjectInfo{Key: key}
|
||||
}
|
||||
}
|
||||
|
||||
actions := make([]RestoreAction, 0, len(candidates))
|
||||
for key, obj := range candidates {
|
||||
rel, include, err := restoreLocalRelativePathForKey(prefix, normalizeRemoteKey(current.CurrentManifestKey), key, opts.IncludeAudio)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("map remote key %q: %w", key, err)
|
||||
}
|
||||
if !include {
|
||||
continue
|
||||
}
|
||||
localPath, err := joinWithinSessionRoot(sessionPaths.Root, rel)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("map remote key %q: %w", key, err)
|
||||
}
|
||||
|
||||
action, err := classifyRestoreAction(ctx, store, obj, rel, localPath, opts.Force)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("classify remote key %q: %w", key, err)
|
||||
}
|
||||
actions = append(actions, action)
|
||||
}
|
||||
|
||||
sort.Slice(actions, func(i, j int) bool {
|
||||
if actions[i].LocalRelativePath == actions[j].LocalRelativePath {
|
||||
return actions[i].RemoteKey < actions[j].RemoteKey
|
||||
}
|
||||
return actions[i].LocalRelativePath < actions[j].LocalRelativePath
|
||||
})
|
||||
|
||||
plan := &RestorePlan{Actions: actions}
|
||||
for _, action := range actions {
|
||||
switch action.Kind {
|
||||
case RestoreActionDownload:
|
||||
plan.DownloadCount++
|
||||
case RestoreActionSkipSame:
|
||||
plan.SkipSameCount++
|
||||
case RestoreActionConflict:
|
||||
plan.ConflictCount++
|
||||
}
|
||||
}
|
||||
|
||||
_ = opts.DryRun
|
||||
return plan, nil
|
||||
}
|
||||
|
||||
func normalizeRemoteKey(v string) string {
|
||||
return strings.Trim(strings.ReplaceAll(strings.TrimSpace(v), "\\", "/"), "/")
|
||||
}
|
||||
|
||||
func restoreLocalRelativePathForKey(sessionPrefix, currentManifestKey, key string, includeAudio bool) (string, bool, error) {
|
||||
if key == "" {
|
||||
return "", false, nil
|
||||
}
|
||||
if key == currentManifestKey {
|
||||
return config.PathManifestFile, true, nil
|
||||
}
|
||||
if !strings.HasPrefix(key, sessionPrefix) {
|
||||
return "", false, fmt.Errorf("key is outside resolved session prefix %q", sessionPrefix)
|
||||
}
|
||||
|
||||
rel := strings.TrimPrefix(key, sessionPrefix)
|
||||
rel = strings.TrimSpace(rel)
|
||||
if rel == "" {
|
||||
return "", false, nil
|
||||
}
|
||||
|
||||
cleanRel := path.Clean(rel)
|
||||
if cleanRel == "." || cleanRel == "" {
|
||||
return "", false, nil
|
||||
}
|
||||
if cleanRel == ".." || strings.HasPrefix(cleanRel, "../") || strings.HasPrefix(cleanRel, "/") {
|
||||
return "", false, fmt.Errorf("key relative path %q escapes session scope", rel)
|
||||
}
|
||||
|
||||
if cleanRel == config.PathManifestFile {
|
||||
return config.PathManifestFile, true, nil
|
||||
}
|
||||
if strings.HasPrefix(cleanRel, config.S3CurrentSegment+"/") {
|
||||
return "", false, nil
|
||||
}
|
||||
if strings.HasPrefix(cleanRel, config.S3RunsSegment+"/") {
|
||||
return "", false, nil
|
||||
}
|
||||
|
||||
excludedRoots := []string{
|
||||
config.PathLogsDirSegment,
|
||||
config.PathReportsDirSegment,
|
||||
config.PathConfigDirSegment,
|
||||
config.PathInputsDirSegment,
|
||||
}
|
||||
for _, root := range excludedRoots {
|
||||
if cleanRel == root || strings.HasPrefix(cleanRel, root+"/") {
|
||||
return "", false, nil
|
||||
}
|
||||
}
|
||||
|
||||
if cleanRel == config.PathTranscriptsSegment || strings.HasPrefix(cleanRel, config.PathTranscriptsSegment+"/") {
|
||||
return cleanRel, true, nil
|
||||
}
|
||||
if cleanRel == config.PathArtifactsDirSegment || strings.HasPrefix(cleanRel, config.PathArtifactsDirSegment+"/") {
|
||||
return cleanRel, true, nil
|
||||
}
|
||||
if includeAudio && (cleanRel == config.PathAudioDirSegment || strings.HasPrefix(cleanRel, config.PathAudioDirSegment+"/")) {
|
||||
return cleanRel, true, nil
|
||||
}
|
||||
|
||||
return "", false, nil
|
||||
}
|
||||
|
||||
func joinWithinSessionRoot(sessionRoot, relative string) (string, error) {
|
||||
if strings.TrimSpace(sessionRoot) == "" {
|
||||
return "", fmt.Errorf("session root is required")
|
||||
}
|
||||
cleanRel := path.Clean(strings.TrimSpace(relative))
|
||||
if cleanRel == "." || cleanRel == "" {
|
||||
return "", fmt.Errorf("relative path is required")
|
||||
}
|
||||
if cleanRel == ".." || strings.HasPrefix(cleanRel, "../") || strings.HasPrefix(cleanRel, "/") {
|
||||
return "", fmt.Errorf("relative path escapes session root")
|
||||
}
|
||||
abs := filepath.Clean(filepath.Join(sessionRoot, filepath.FromSlash(cleanRel)))
|
||||
root := filepath.Clean(sessionRoot)
|
||||
if abs != root && !strings.HasPrefix(abs, root+string(filepath.Separator)) {
|
||||
return "", fmt.Errorf("resolved local path escapes session root")
|
||||
}
|
||||
return abs, nil
|
||||
}
|
||||
|
||||
func classifyRestoreAction(
|
||||
ctx context.Context,
|
||||
store storage.ObjectStore,
|
||||
object storage.ObjectInfo,
|
||||
localRelPath string,
|
||||
localPath string,
|
||||
force bool,
|
||||
) (RestoreAction, error) {
|
||||
action := RestoreAction{
|
||||
RemoteKey: normalizeRemoteKey(object.Key),
|
||||
LocalRelativePath: localRelPath,
|
||||
LocalPath: localPath,
|
||||
Size: object.Size,
|
||||
ETag: object.ETag,
|
||||
}
|
||||
|
||||
info, err := os.Stat(localPath)
|
||||
if err != nil {
|
||||
if os.IsNotExist(err) {
|
||||
action.Kind = RestoreActionDownload
|
||||
action.Reason = "local file missing"
|
||||
return action, nil
|
||||
}
|
||||
return RestoreAction{}, fmt.Errorf("stat local file: %w", err)
|
||||
}
|
||||
|
||||
action.ExistsLocal = true
|
||||
if info.IsDir() {
|
||||
action.Kind = RestoreActionConflict
|
||||
action.Conflict = true
|
||||
action.Reason = "local path is a directory"
|
||||
return action, nil
|
||||
}
|
||||
|
||||
if object.Size > 0 && info.Size() != object.Size {
|
||||
if force {
|
||||
action.Kind = RestoreActionDownload
|
||||
action.Reason = "local file differs (size mismatch); overwrite with --force"
|
||||
return action, nil
|
||||
}
|
||||
action.Kind = RestoreActionConflict
|
||||
action.Conflict = true
|
||||
action.Reason = "local file differs (size mismatch)"
|
||||
return action, nil
|
||||
}
|
||||
|
||||
localDigest, err := artifacts.SHA256File(localPath)
|
||||
if err != nil {
|
||||
return RestoreAction{}, fmt.Errorf("checksum local file: %w", err)
|
||||
}
|
||||
remotePath, err := downloadObjectToTemp(ctx, store, action.RemoteKey, "narratio-restore-plan-remote-*.tmp")
|
||||
if err != nil {
|
||||
return RestoreAction{}, fmt.Errorf("download remote object: %w", err)
|
||||
}
|
||||
defer func() { _ = os.Remove(remotePath) }()
|
||||
|
||||
remoteDigest, err := artifacts.SHA256File(remotePath)
|
||||
if err != nil {
|
||||
return RestoreAction{}, fmt.Errorf("checksum remote object: %w", err)
|
||||
}
|
||||
|
||||
if remoteDigest == localDigest {
|
||||
action.Kind = RestoreActionSkipSame
|
||||
action.SameLocal = true
|
||||
action.Reason = "local file matches remote content"
|
||||
return action, nil
|
||||
}
|
||||
|
||||
if force {
|
||||
action.Kind = RestoreActionDownload
|
||||
action.Reason = "local file differs; overwrite with --force"
|
||||
return action, nil
|
||||
}
|
||||
|
||||
action.Kind = RestoreActionConflict
|
||||
action.Conflict = true
|
||||
action.Reason = "local file differs"
|
||||
return action, nil
|
||||
}
|
||||
|
||||
func writeRestorePlan(out io.Writer, current *RemoteCurrentState, plan *RestorePlan, opts RestorePlanOptions) error {
|
||||
if out == nil {
|
||||
return fmt.Errorf("output writer is required")
|
||||
}
|
||||
if current == nil {
|
||||
return fmt.Errorf("remote current state is required")
|
||||
}
|
||||
if plan == nil {
|
||||
return fmt.Errorf("restore plan is required")
|
||||
}
|
||||
|
||||
if _, err := fmt.Fprintf(
|
||||
out,
|
||||
"restore plan: session %s/%s run=%s actions=%d download=%d skip_same=%d conflict=%d dry_run=%t force=%t include_audio=%t\n",
|
||||
current.Campaign,
|
||||
current.SessionID,
|
||||
current.RunID,
|
||||
len(plan.Actions),
|
||||
plan.DownloadCount,
|
||||
plan.SkipSameCount,
|
||||
plan.ConflictCount,
|
||||
opts.DryRun,
|
||||
opts.Force,
|
||||
opts.IncludeAudio,
|
||||
); err != nil {
|
||||
return err
|
||||
}
|
||||
for _, action := range plan.Actions {
|
||||
if _, err := fmt.Fprintf(out, "%s %s <- %s", action.Kind, action.LocalRelativePath, action.RemoteKey); err != nil {
|
||||
return err
|
||||
}
|
||||
if strings.TrimSpace(action.Reason) != "" {
|
||||
if _, err := fmt.Fprintf(out, " (%s)", action.Reason); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
if _, err := fmt.Fprintln(out); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
Reference in New Issue
Block a user