package stage import ( "context" "fmt" "io/fs" "os" "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/manifest" ) type archiveStage struct{} var archivePrerequisiteStages = []string{ "prepare", "transcribe", "merge", "polish", "normalize", "trim", "analyze", } var archiveRunUploadDirs = []string{ "inputs", "transcripts", "artifacts", "reports", "config", "logs", } func (archiveStage) Name() string { return "archive" } func (archiveStage) Declares() IODecl { return IODecl{ Inputs: []artifacts.Ref{ {Kind: "manifest", Category: "input", RelativePath: "manifest.json"}, }, } } func (archiveStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*StageResult, error) { if env == nil || env.Config == nil || env.Config.Pipeline == nil || env.Config.Session == nil { return nil, fmt.Errorf("archive: resolved config must include pipeline and session") } if archiveDisabled(env) { return &StageResult{ Metadata: map[string]any{ "stage": "archive", "skipped": true, "archive_enabled": false, "audio_upload_skipped": true, }, }, nil } if archiveRunUploadDisabled(env) { return &StageResult{ Metadata: map[string]any{ "stage": "archive", "skipped": true, "upload_run_enabled": false, "audio_upload_skipped": true, }, }, nil } if err := validateArchivePrerequisites(m); err != nil { return nil, fmt.Errorf("archive: %w", err) } if env.ObjectStore == nil { return nil, fmt.Errorf("archive: remote object store backend is required when archive run upload is enabled") } workDir, err := archiveWorkDir(env, m) if err != nil { return nil, fmt.Errorf("archive: resolve local workdir: %w", err) } info, err := os.Stat(workDir) if err != nil { return nil, fmt.Errorf("archive: local workdir %q: %w", workDir, err) } if !info.IsDir() { return nil, fmt.Errorf("archive: local workdir %q is not a directory", workDir) } runPrefix, err := archiveRunPrefix(env, m) if err != nil { return nil, fmt.Errorf("archive: resolve s3 run prefix: %w", err) } bucket := archiveBucket(env, m) if bucket == "" { return nil, fmt.Errorf("archive: resolve s3 bucket: bucket is required") } relFiles, err := collectArchiveRunFiles(workDir) if err != nil { return nil, fmt.Errorf("archive: collect run files: %w", err) } uploaded := make([]string, 0, len(relFiles)) for _, rel := range relFiles { localPath := filepath.Join(workDir, filepath.FromSlash(rel)) key := artifacts.S3RunRelativeDestinationKey(runPrefix, rel) if _, err := env.ObjectStore.Upload(ctx, localPath, key, storage.UploadOptions{}); err != nil { return nil, fmt.Errorf("archive: upload %q to %q: %w", rel, key, err) } uploaded = append(uploaded, rel) } return &StageResult{ Metadata: map[string]any{ "stage": "archive", "uploaded": true, "s3_bucket": bucket, "s3_run_prefix": runPrefix, "files_uploaded": len(uploaded), "uploaded_paths": uploaded, "audio_upload_skipped": true, }, }, nil } func archiveDisabled(env *Env) bool { cfg := env.Config.Pipeline.Archive if cfg == nil { return true } return cfg.Enabled != nil && !*cfg.Enabled } func archiveRunUploadDisabled(env *Env) bool { cfg := env.Config.Pipeline.Archive if cfg == nil { return true } return cfg.UploadRun != nil && !*cfg.UploadRun } func validateArchivePrerequisites(m *manifest.Manifest) error { if m == nil { return fmt.Errorf("manifest is required") } for _, stageName := range archivePrerequisiteStages { sr := m.Stages[stageName] if sr == nil { return fmt.Errorf("prerequisite stage %q has not succeeded", stageName) } if sr.Status != manifest.StatusSucceeded { return fmt.Errorf("prerequisite stage %q status is %q (want %q)", stageName, sr.Status, manifest.StatusSucceeded) } } return nil } func archiveWorkDir(env *Env, m *manifest.Manifest) (string, error) { workDir := strings.TrimSpace(m.LocalWorkDir) if workDir != "" { cleaned := filepath.Clean(workDir) if info, err := os.Stat(cleaned); err == nil && info.IsDir() { return cleaned, nil } } sessionID := strings.TrimSpace(env.Config.Session.SessionID) if sessionID == "" { sessionID = strings.TrimSpace(m.SessionID) } campaign := strings.TrimSpace(env.Config.Session.Campaign) if campaign == "" { campaign = strings.TrimSpace(m.Campaign) } runID := strings.TrimSpace(m.RunID) if runID == "" { return "", fmt.Errorf("run id is required") } if campaign == "" || sessionID == "" { return "", fmt.Errorf("campaign and session id are required") } runScoped := artifacts.SessionRunWorkDir(env.Config.Pipeline.Workspace.Root, campaign, sessionID, runID) if info, err := os.Stat(runScoped); err == nil && info.IsDir() { return runScoped, nil } legacy := artifacts.SessionWorkDir(env.Config.Pipeline.Workspace.Root, sessionID) return legacy, nil } func archiveRunPrefix(env *Env, m *manifest.Manifest) (string, error) { runPrefix := strings.TrimSpace(m.S3RunPrefix) if runPrefix != "" { return runPrefix, nil } sessionID := strings.TrimSpace(env.Config.Session.SessionID) if sessionID == "" { sessionID = strings.TrimSpace(m.SessionID) } campaign := strings.TrimSpace(env.Config.Session.Campaign) if campaign == "" { campaign = strings.TrimSpace(m.Campaign) } runID := strings.TrimSpace(m.RunID) if runID == "" { return "", fmt.Errorf("run id is required") } if env.Config.Pipeline.Storage.S3 == nil { return "", fmt.Errorf("pipeline.storage.s3 configuration is required") } sessionPrefix := artifacts.S3SessionPrefix(env.Config.Pipeline.Storage.S3.RootPrefix, campaign, sessionID) if strings.TrimSpace(sessionPrefix) == "" { return "", fmt.Errorf("session prefix is required") } return artifacts.S3RunPrefix(sessionPrefix, runID), nil } func archiveBucket(env *Env, m *manifest.Manifest) string { if m != nil && strings.TrimSpace(m.S3Bucket) != "" { return strings.TrimSpace(m.S3Bucket) } if env == nil || env.Config == nil || env.Config.Pipeline == nil || env.Config.Pipeline.Storage.S3 == nil { return "" } return strings.TrimSpace(env.Config.Pipeline.Storage.S3.Bucket) } func collectArchiveRunFiles(workDir string) ([]string, error) { files := make([]string, 0, 64) for _, dirName := range archiveRunUploadDirs { fullDir := filepath.Join(workDir, dirName) info, err := os.Stat(fullDir) if err != nil { if os.IsNotExist(err) { continue } return nil, fmt.Errorf("stat %q: %w", fullDir, err) } if !info.IsDir() { continue } if err := filepath.WalkDir(fullDir, func(path string, d fs.DirEntry, walkErr error) error { if walkErr != nil { return walkErr } if d.IsDir() { return nil } rel, err := filepath.Rel(workDir, path) if err != nil { return fmt.Errorf("relative path from %q to %q: %w", workDir, path, err) } rel = filepath.ToSlash(rel) files = append(files, rel) return nil }); err != nil { return nil, fmt.Errorf("walk %q: %w", fullDir, err) } } manifestPath := filepath.Join(workDir, "manifest.json") manifestInfo, err := os.Stat(manifestPath) if err != nil { if os.IsNotExist(err) { return nil, fmt.Errorf("manifest.json not found in workdir %q", workDir) } return nil, fmt.Errorf("stat %q: %w", manifestPath, err) } if manifestInfo.IsDir() { return nil, fmt.Errorf("manifest path %q is a directory", manifestPath) } files = append(files, "manifest.json") sort.Strings(files) return files, nil }