9 Commits

35 changed files with 1276 additions and 585 deletions

View File

@@ -2,7 +2,7 @@
Narratio is a stage-driven Go orchestrator for turning D&D session audio into polished transcripts and generated artifacts. Narratio is a stage-driven Go orchestrator for turning D&D session audio into polished transcripts and generated artifacts.
It runs a deterministic workflow across `prepare`, `transcribe`, `merge`, `polish`, `normalize`, `trim`, `analyze`, and `publish`, with manifest-driven resume and restore support. It runs a deterministic workflow across `prepare`, `transcribe`, `merge`, `polish`, `normalize`, `trim`, `analyze`, and `publish`, with manifest-driven continuation and restore support.
```bash ```bash
narratio run 2026-04-04 narratio run 2026-04-04

View File

@@ -13,7 +13,6 @@ This runs the canonical full pipeline for session `2026-04-04`.
Top-level commands: Top-level commands:
- `run <session_id>`: run full stage order. - `run <session_id>`: run full stage order.
- `resume <session_id>`: continue from first non-succeeded stage.
- `run-stage <stage> <session_id>`: run one stage. - `run-stage <stage> <session_id>`: run one stage.
- `analyze <session_id>`: force-run analyze. - `analyze <session_id>`: force-run analyze.
- `publish <session_id>`: force-run publish. - `publish <session_id>`: force-run publish.
@@ -77,19 +76,9 @@ Behavior:
- evaluates full stage order; - evaluates full stage order;
- skips already-succeeded stages unless `--force` is set; - skips already-succeeded stages unless `--force` is set;
- continues interrupted or partially completed sessions by running non-succeeded stages;
- writes session and run manifests. - writes session and run manifests.
### `resume`
```bash
narratio resume <session_id> [--force] [--artifacts <name[,name...]>] [...common config flags]
```
Behavior:
- when not forced, starts at first non-succeeded stage in manifest order;
- with `--force`, reevaluates the selected stage list as runnable.
### `run-stage` ### `run-stage`
```bash ```bash
@@ -253,7 +242,7 @@ Behavior:
## `--artifacts` Selection Rules ## `--artifacts` Selection Rules
- accepted on `run`, `resume`, `run-stage`, `analyze`, and `publish`; - accepted on `run`, `run-stage`, `analyze`, and `publish`;
- names must exist in `pipeline.scriptorium.artifacts`; - names must exist in `pipeline.scriptorium.artifacts`;
- empty entries are invalid; - empty entries are invalid;
- repeated names are deduplicated. - repeated names are deduplicated.

View File

@@ -4,33 +4,6 @@ Operator workflow for running, recovering, and publishing Narratio sessions.
For command syntax, see [docs/cli.md](./cli.md). For field-level config, see [docs/config.md](./config.md). For command syntax, see [docs/cli.md](./cli.md). For field-level config, see [docs/config.md](./config.md).
## Standard Session Workflow
1. Select pipeline/campaign/session config.
2. Validate session readiness:
```bash
narratio session validate 2026-04-04
```
3. (Optional) inspect stage decisions:
```bash
narratio session plan 2026-04-04
```
4. Run the pipeline:
```bash
narratio run 2026-04-04
```
5. Check state:
```bash
narratio session status 2026-04-04
```
## Campaign and Session Selection ## Campaign and Session Selection
Campaign selection priority: Campaign selection priority:
@@ -63,7 +36,34 @@ narratio session init 2026-04-04 --remote --force
If `campaign.yml` sets `session_template_file`, `session init` renders it. Template variables must resolve to concrete values. If `campaign.yml` sets `session_template_file`, `session init` renders it. Template variables must resolve to concrete values.
## Stage Execution and Resume Behavior ## Standard Session Workflow
1. Select pipeline/campaign/session config.
2. Validate session readiness:
```bash
narratio session validate 2026-04-04
```
3. (Optional) inspect stage decisions:
```bash
narratio session plan 2026-04-04
```
4. Run the pipeline:
```bash
narratio run 2026-04-04
```
5. Check state:
```bash
narratio session status 2026-04-04
```
## Stage Execution and Continuation Behavior
Canonical stage order: Canonical stage order:
@@ -80,7 +80,7 @@ Canonical stage order:
Execution rules: Execution rules:
- succeeded stages are skipped unless `--force` is set; - succeeded stages are skipped unless `--force` is set;
- `resume` starts at first non-succeeded stage; - `run` continues interrupted or partially completed sessions by running non-succeeded stages;
- force rerunning a succeeded upstream stage marks succeeded downstream stages as `stale`. - force rerunning a succeeded upstream stage marks succeeded downstream stages as `stale`.
Single-stage execution: Single-stage execution:
@@ -91,7 +91,7 @@ narratio run-stage normalize 2026-04-04 --force
## Artifact Selection ## Artifact Selection
`--artifacts` can be used on `run`, `resume`, `run-stage`, `analyze`, and `publish`. `--artifacts` can be used on `run`, `run-stage`, `analyze`, and `publish`.
Selection behavior: Selection behavior:

View File

@@ -6,7 +6,7 @@ Canonical contributor workflow and engineering conventions for implemented Narra
## Repository layout ## Repository layout
- `cmd/narratio/`: CLI entrypoint. - `cmd/narratio/`: CLI entrypoint.
- `internal/app/`: command handlers, plan/run/resume orchestration, cleanup gates, secrets loading. - `internal/app/`: command handlers, run/stage orchestration, cleanup gates, secrets loading.
- `internal/config/`: strict YAML loading, defaults, and validation. - `internal/config/`: strict YAML loading, defaults, and validation.
- `internal/stage/`: stage implementations and stage registry/order. - `internal/stage/`: stage implementations and stage registry/order.
- `internal/adapters/`: external boundary adapters (WhisperX, Seriatim, Audita, Scriptorium, storage, notify). - `internal/adapters/`: external boundary adapters (WhisperX, Seriatim, Audita, Scriptorium, storage, notify).

View File

@@ -120,7 +120,7 @@ func TestRunStageArtifactsDoesNotImplyForce(t *testing.T) {
} }
} }
func TestResumeArtifactsWithSucceededAnalyzeSkipsUnlessForced(t *testing.T) { func TestRunArtifactsWithSucceededAnalyzeSkipsUnlessForced(t *testing.T) {
workspaceRoot := t.TempDir() workspaceRoot := t.TempDir()
pipelinePath, campaignPath, sessionPath := writeValidConfigFilesWithScriptoriumArtifacts(t, workspaceRoot) pipelinePath, campaignPath, sessionPath := writeValidConfigFilesWithScriptoriumArtifacts(t, workspaceRoot)
manifestPath := filepath.Join(workspaceRoot, "work", "sample-campaign", "2026-05-03", "manifest.json") manifestPath := filepath.Join(workspaceRoot, "work", "sample-campaign", "2026-05-03", "manifest.json")
@@ -135,16 +135,16 @@ func TestResumeArtifactsWithSucceededAnalyzeSkipsUnlessForced(t *testing.T) {
} }
var out bytes.Buffer var out bytes.Buffer
err := Resume( err := Run(
context.Background(), context.Background(),
[]string{"2026-05-03", "--config", pipelinePath, "--campaign-file", campaignPath, "--session", sessionPath, "--artifacts", "session_recap"}, []string{"2026-05-03", "--config", pipelinePath, "--campaign-file", campaignPath, "--session", sessionPath, "--artifacts", "session_recap"},
&out, &out,
) )
if err != nil { if err != nil {
t.Fatalf("Resume() error = %v", err) t.Fatalf("Run() error = %v", err)
} }
if !strings.Contains(out.String(), "has no remaining stages") { if !strings.Contains(out.String(), "executed=0 skipped=9") {
t.Fatalf("output = %q, want no remaining stages", out.String()) t.Fatalf("output = %q, want all stages skipped", out.String())
} }
} }

View File

@@ -7,7 +7,7 @@ import (
"strings" "strings"
) )
var supportedCommands = []string{"run", "run-stage", "resume", "analyze", "publish", "clean", "session"} var supportedCommands = []string{"run", "run-stage", "analyze", "publish", "clean", "session"}
// Execute dispatches CLI commands and returns a process exit code. // Execute dispatches CLI commands and returns a process exit code.
func Execute(args []string, stdout, stderr io.Writer) int { func Execute(args []string, stdout, stderr io.Writer) int {
@@ -24,8 +24,6 @@ func Execute(args []string, stdout, stderr io.Writer) int {
switch cmd { switch cmd {
case "run": case "run":
err = Run(ctx, cmdArgs, stdout) err = Run(ctx, cmdArgs, stdout)
case "resume":
err = Resume(ctx, cmdArgs, stdout)
case "run-stage": case "run-stage":
err = RunStage(ctx, cmdArgs, stdout) err = RunStage(ctx, cmdArgs, stdout)
case "analyze": case "analyze":

View File

@@ -34,7 +34,6 @@ func TestExecuteValidCommands(t *testing.T) {
{name: "run", args: []string{"run", "2026-05-03", "--config", pipelinePath, "--campaign-file", campaignPath, "--session", sessionPath}, wantOut: "narratio run: session 2026-05-03; executed=9 skipped=0; manifest="}, {name: "run", args: []string{"run", "2026-05-03", "--config", pipelinePath, "--campaign-file", campaignPath, "--session", sessionPath}, wantOut: "narratio run: session 2026-05-03; executed=9 skipped=0; manifest="},
{name: "session plan", args: []string{"session", "plan", "2026-05-03", "--config", pipelinePath, "--campaign-file", campaignPath, "--session", sessionPath}, wantOut: "prepare: skip\ntranscribe: skip\nmerge: skip\npolish: skip\nnormalize: skip\ntrim: skip\nanalyze: skip\npublish: skip\nnotify: skip"}, {name: "session plan", args: []string{"session", "plan", "2026-05-03", "--config", pipelinePath, "--campaign-file", campaignPath, "--session", sessionPath}, wantOut: "prepare: skip\ntranscribe: skip\nmerge: skip\npolish: skip\nnormalize: skip\ntrim: skip\nanalyze: skip\npublish: skip\nnotify: skip"},
{name: "session status", args: []string{"session", "status", "2026-05-03", "--config", pipelinePath, "--campaign-file", campaignPath, "--session", sessionPath}, wantOut: "Session: 2026-05-03"}, {name: "session status", args: []string{"session", "status", "2026-05-03", "--config", pipelinePath, "--campaign-file", campaignPath, "--session", sessionPath}, wantOut: "Session: 2026-05-03"},
{name: "resume", args: []string{"resume", "2026-05-03", "--config", pipelinePath, "--campaign-file", campaignPath, "--session", sessionPath}, wantOut: "narratio resume: session 2026-05-03 has no remaining stages"},
{name: "run-stage", args: []string{"run-stage", "polish", "2026-05-03", "--config", pipelinePath, "--campaign-file", campaignPath, "--session", sessionPath}, wantOut: "narratio run-stage: stage=polish executed=0 skipped=1 force=false; manifest="}, {name: "run-stage", args: []string{"run-stage", "polish", "2026-05-03", "--config", pipelinePath, "--campaign-file", campaignPath, "--session", sessionPath}, wantOut: "narratio run-stage: stage=polish executed=0 skipped=1 force=false; manifest="},
} }
@@ -66,7 +65,7 @@ func TestExecuteMissingRequiredFlags(t *testing.T) {
{name: "run missing session", args: []string{"run"}, want: "run: session_id is required"}, {name: "run missing session", args: []string{"run"}, want: "run: session_id is required"},
{name: "plan old top-level removed", args: []string{"plan"}, want: `unknown command: "plan"`}, {name: "plan old top-level removed", args: []string{"plan"}, want: `unknown command: "plan"`},
{name: "status old top-level removed", args: []string{"status"}, want: `unknown command: "status"`}, {name: "status old top-level removed", args: []string{"status"}, want: `unknown command: "status"`},
{name: "resume missing session", args: []string{"resume"}, want: "resume: session_id is required"}, {name: "resume removed", args: []string{"resume"}, want: `unknown command: "resume"`},
{name: "run-stage missing name", args: []string{"run-stage", "--config", "a", "--session", "b"}, want: "run-stage: expected stage name and session_id"}, {name: "run-stage missing name", args: []string{"run-stage", "--config", "a", "--session", "b"}, want: "run-stage: expected stage name and session_id"},
{name: "run-stage missing session", args: []string{"run-stage", "polish"}, want: "run-stage: expected stage name and session_id"}, {name: "run-stage missing session", args: []string{"run-stage", "polish"}, want: "run-stage: expected stage name and session_id"},
{name: "run missing config uses defaults", args: []string{"run", "2026-05-03", "--session", "session.yml"}, want: "run: no pipeline config path provided and no default pipeline config found; searched:"}, {name: "run missing config uses defaults", args: []string{"run", "2026-05-03", "--session", "session.yml"}, want: "run: no pipeline config path provided and no default pipeline config found; searched:"},

View File

@@ -10,7 +10,6 @@ import (
"strings" "strings"
"gitea.maximumdirect.net/eric/narratio/internal/adapters/storage" "gitea.maximumdirect.net/eric/narratio/internal/adapters/storage"
"gitea.maximumdirect.net/eric/narratio/internal/artifacts"
"gitea.maximumdirect.net/eric/narratio/internal/config" "gitea.maximumdirect.net/eric/narratio/internal/config"
"gitea.maximumdirect.net/eric/narratio/internal/manifest" "gitea.maximumdirect.net/eric/narratio/internal/manifest"
) )
@@ -64,26 +63,18 @@ func sessionSourceSummary(cfg *config.Config) string {
} }
func validateStableInputFindings(cfg *config.Config) []finding { func validateStableInputFindings(cfg *config.Config) []finding {
items := []struct { checks := inspectStableInputs(cfg)
name string out := make([]finding, 0, len(checks))
in config.ResolvedInputFile for _, check := range checks {
}{ if check.Err != nil {
{"speakers", cfg.StableInputs.SpeakersFile}, msg := check.Name + ": " + check.Err.Error()
{"autocorrect", cfg.StableInputs.AutocorrectFile}, if strings.TrimSpace(check.Path) != "" {
{"glossary", cfg.StableInputs.GlossaryFile}, msg = fmt.Sprintf("%s missing: %v", check.Name, check.Err)
} }
out := make([]finding, 0, len(items)) out = append(out, errorFinding("inputs", msg))
for _, item := range items {
path, err := resolveHelperConfigRelativePath(item.in)
if err != nil {
out = append(out, errorFinding("inputs", item.name+": "+err.Error()))
continue continue
} }
if _, err := os.Stat(path); err != nil { out = append(out, okFinding("inputs", check.Name+": "+check.Path))
out = append(out, errorFinding("inputs", fmt.Sprintf("%s missing: %v", item.name, err)))
} else {
out = append(out, okFinding("inputs", item.name+": "+path))
}
} }
return out return out
} }
@@ -103,76 +94,22 @@ func resolveHelperConfigRelativePath(input config.ResolvedInputFile) (string, er
} }
func validateLocalAudioFindings(cfg *config.Config) []finding { func validateLocalAudioFindings(cfg *config.Config) []finding {
if cfg.Session.Inputs.AudioS3 != nil { check := inspectLocalAudioPresence(cfg)
if !check.Checked {
return nil return nil
} }
audioDir := strings.TrimSpace(cfg.Session.Inputs.AudioDir) if check.Err != nil {
if audioDir == "" && len(cfg.Session.Inputs.AudioFiles) == 0 { return []finding{errorFinding("audio", check.Err.Error())}
return []finding{errorFinding("audio", "audio_dir, audio_files, or audio_s3 is required")}
} }
base := filepath.Dir(cfg.SessionPath) return []finding{okFinding("audio", fmt.Sprintf("%d local audio file(s)", len(check.Paths)))}
paths := []string{}
if audioDir != "" {
dir := audioDir
if !filepath.IsAbs(dir) {
dir = filepath.Join(base, dir)
}
matches, err := filepath.Glob(filepath.Join(dir, "*.flac"))
if err != nil || len(matches) == 0 {
return []finding{errorFinding("audio", "no .flac files found in "+dir)}
}
paths = append(paths, matches...)
}
for _, file := range cfg.Session.Inputs.AudioFiles {
p := file
if !filepath.IsAbs(p) {
p = filepath.Join(base, p)
}
paths = append(paths, p)
}
for _, p := range paths {
if _, err := os.Stat(p); err != nil {
return []finding{errorFinding("audio", fmt.Sprintf("audio file missing: %v", err))}
}
}
return []finding{okFinding("audio", fmt.Sprintf("%d local audio file(s)", len(paths)))}
} }
func validateRemoteAudioFinding(ctx context.Context, cfg *config.Config, store storage.ObjectStore) finding { func validateRemoteAudioFinding(ctx context.Context, cfg *config.Config, store storage.ObjectStore) finding {
sessionPrefix := artifacts.S3SessionPrefix(cfg.Pipeline.Storage.S3.RootPrefix, cfg.Session.Campaign, cfg.Session.SessionID) check := inspectRemoteAudioPresence(ctx, cfg, store)
audioPrefix := artifacts.S3AudioPrefix(sessionPrefix, cfg.Session.Inputs.AudioS3.Prefix) if check.Err != nil {
objects, err := store.List(ctx, audioPrefix) return errorFinding("audio", check.Err.Error())
if err != nil {
return errorFinding("audio", err.Error())
} }
count := 0 return okFinding("audio", fmt.Sprintf("%d remote .flac object(s)", len(check.Keys)))
for _, obj := range objects {
if strings.HasSuffix(strings.ToLower(obj.Key), ".flac") {
count++
}
}
if count == 0 {
return errorFinding("audio", "no remote .flac objects found under "+audioPrefix)
}
return okFinding("audio", fmt.Sprintf("%d remote .flac object(s)", count))
}
func validatePreviousArtifactFindings(ctx context.Context, cfg *config.Config, store storage.ObjectStore, requirements []artifacts.PreviousArtifactRequirement) []finding {
out := []finding{}
prefix := artifacts.S3SessionPrefix(cfg.Pipeline.Storage.S3.RootPrefix, cfg.Session.Campaign, cfg.Session.PreviousSessionID)
_, err := artifacts.LoadCurrentState(ctx, store, prefix, artifacts.CurrentStateValidation{
ExpectedSessionID: strings.TrimSpace(cfg.Session.PreviousSessionID),
ExpectedCampaign: strings.TrimSpace(cfg.Session.Campaign),
ValidateRunID: true,
})
if err != nil {
out = append(out, errorFinding("previous", fmt.Sprintf("remote %v", err)))
return out
}
for _, req := range requirements {
out = append(out, okFinding("previous", fmt.Sprintf("%s required=%t", req.Name, req.Required)))
}
return out
} }
func loadLocalManifest(ctx context.Context, path string) (*manifest.Manifest, error) { func loadLocalManifest(ctx context.Context, path string) (*manifest.Manifest, error) {

View File

@@ -593,6 +593,53 @@ func TestExecuteLocksRequireSessionID(t *testing.T) {
} }
} }
func TestExecuteLocksMutationRejectsSessionIDMismatch(t *testing.T) {
workspaceRoot := t.TempDir()
pipelinePath, campaignPath, sessionPath := writeValidConfigFiles(t, workspaceRoot)
fake := &storage.FakeBackend{}
var storeInitCalls int
restoreAppConfigTestGlobals(t, fake, &storeInitCalls, []string{sessionPath})
tests := []struct {
name string
args []string
}{
{
name: "add mismatch",
args: []string{
"session", "locks", "add", "2026-05-03", "narratio.transcript.final_trimmed",
"--session-id", "2026-05-04",
"--config", pipelinePath,
"--campaign-file", campaignPath,
"--session", sessionPath,
},
},
{
name: "remove mismatch",
args: []string{
"session", "locks", "remove", "2026-05-03", "narratio.transcript.final_trimmed",
"--session-id", "2026-05-04",
"--config", pipelinePath,
"--campaign-file", campaignPath,
"--session", sessionPath,
},
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
var stdout bytes.Buffer
var stderr bytes.Buffer
code := Execute(tt.args, &stdout, &stderr)
if code == 0 {
t.Fatal("exit code = 0, want non-zero")
}
if !strings.Contains(stderr.String(), "does not match expected session id") {
t.Fatalf("stderr = %q, want session-id mismatch guidance", stderr.String())
}
})
}
}
func TestExecuteLocksCannotModifyStaticLocks(t *testing.T) { func TestExecuteLocksCannotModifyStaticLocks(t *testing.T) {
workspaceRoot := t.TempDir() workspaceRoot := t.TempDir()
pipelinePath, campaignPath, sessionPath := writeValidConfigFiles(t, workspaceRoot) pipelinePath, campaignPath, sessionPath := writeValidConfigFiles(t, workspaceRoot)
@@ -884,6 +931,31 @@ func TestExecuteStatusReportsMissingRemoteCurrentStateWithoutFailing(t *testing.
} }
} }
func TestExecuteStatusReportsPreviousStateReadinessWithoutFailing(t *testing.T) {
workspaceRoot := t.TempDir()
pipelinePath, campaignPath, sessionPath := writeValidConfigFilesWithScriptoriumArtifacts(t, workspaceRoot)
replaceInFileOrFatal(t, pipelinePath, "source: narratio.artifact.session_recap", "source: narratio.previous_session.artifact.session_recap")
replaceInFileOrFatal(t, sessionPath, "session_id: 2026-05-03\n", "session_id: 2026-05-03\nprevious_session_id: 2026-04-26\n")
fake := &storage.FakeBackend{}
var storeInitCalls int
restoreAppConfigTestGlobals(t, fake, &storeInitCalls, []string{sessionPath})
var stdout bytes.Buffer
var stderr bytes.Buffer
code := Execute([]string{
"session", "status", "2026-05-03",
"--config", pipelinePath,
"--campaign-file", campaignPath,
"--session", sessionPath,
}, &stdout, &stderr)
if code != 0 {
t.Fatalf("exit code = %d, want 0; stderr=%q", code, stderr.String())
}
if !strings.Contains(stdout.String(), "Previous-session artifacts: unavailable: remote current run pointer missing") {
t.Fatalf("stdout = %q, want previous readiness unavailable line", stdout.String())
}
}
func TestExecuteSessionValidateReportsPreviousStateFindingAndReturnsFindingError(t *testing.T) { func TestExecuteSessionValidateReportsPreviousStateFindingAndReturnsFindingError(t *testing.T) {
workspaceRoot := t.TempDir() workspaceRoot := t.TempDir()
pipelinePath, campaignPath, sessionPath := writeValidConfigFilesWithScriptoriumArtifacts(t, workspaceRoot) pipelinePath, campaignPath, sessionPath := writeValidConfigFilesWithScriptoriumArtifacts(t, workspaceRoot)

View File

@@ -0,0 +1,274 @@
package app
import (
"context"
"fmt"
"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"
)
type stableInputCheck struct {
Name string
Path string
Err error
}
type localAudioCheck struct {
Checked bool
Paths []string
Err error
}
type remoteAudioCheck struct {
Checked bool
Prefix string
Keys []string
Err error
}
type previousArtifactReadiness struct {
Requirements []artifacts.PreviousArtifactRequirement
MissingID bool
Err error
}
type remoteCurrentStateCheck struct {
State *RemoteCurrentState
Err error
}
type effectiveLocksCheck struct {
Locks *effectiveLocks
Err error
}
func inspectStableInputs(cfg *config.Config) []stableInputCheck {
items := []struct {
name string
in config.ResolvedInputFile
}{
{name: "speakers", in: cfg.StableInputs.SpeakersFile},
{name: "autocorrect", in: cfg.StableInputs.AutocorrectFile},
{name: "glossary", in: cfg.StableInputs.GlossaryFile},
}
out := make([]stableInputCheck, 0, len(items))
for _, item := range items {
path, err := resolveHelperConfigRelativePath(item.in)
if err != nil {
out = append(out, stableInputCheck{Name: item.name, Err: err})
continue
}
if _, err := os.Stat(path); err != nil {
out = append(out, stableInputCheck{Name: item.name, Path: path, Err: err})
continue
}
out = append(out, stableInputCheck{Name: item.name, Path: path})
}
return out
}
func inspectLocalAudioPresence(cfg *config.Config) localAudioCheck {
if cfg.Session.Inputs.AudioS3 != nil {
return localAudioCheck{}
}
sessionDir := filepath.Dir(cfg.SessionPath)
resolved, err := resolveLocalInspectionAudioPaths(sessionDir, cfg.Session.Inputs)
if err != nil {
return localAudioCheck{Checked: true, Err: err}
}
return localAudioCheck{
Checked: true,
Paths: resolved,
}
}
func inspectRemoteAudioPresence(ctx context.Context, cfg *config.Config, store storage.ObjectStore) remoteAudioCheck {
if cfg.Session.Inputs.AudioS3 == nil {
return remoteAudioCheck{}
}
if store == nil {
return remoteAudioCheck{Checked: true, Err: fmt.Errorf("storage backend is required for remote audio checks")}
}
sessionPrefix := artifacts.S3SessionPrefix(cfg.Pipeline.Storage.S3.RootPrefix, cfg.Session.Campaign, cfg.Session.SessionID)
audioPrefix := artifacts.S3AudioPrefix(sessionPrefix, cfg.Session.Inputs.AudioS3.Prefix)
objects, err := store.List(ctx, audioPrefix)
if err != nil {
return remoteAudioCheck{Checked: true, Prefix: audioPrefix, Err: err}
}
keys := make([]string, 0, len(objects))
seenBase := map[string]string{}
for _, obj := range objects {
key := strings.TrimSpace(obj.Key)
if key == "" || strings.HasSuffix(key, "/") || !isInspectionFlacPath(key) {
continue
}
base := path.Base(key)
if prev, exists := seenBase[base]; exists && prev != key {
return remoteAudioCheck{
Checked: true,
Prefix: audioPrefix,
Err: fmt.Errorf("duplicate s3 audio basename %q from %q and %q", base, prev, key),
}
}
seenBase[base] = key
keys = append(keys, key)
}
sort.Strings(keys)
if len(keys) == 0 {
return remoteAudioCheck{
Checked: true,
Prefix: audioPrefix,
Err: fmt.Errorf("no .flac files found under s3 audio prefix %q", audioPrefix),
}
}
return remoteAudioCheck{
Checked: true,
Prefix: audioPrefix,
Keys: keys,
}
}
func inspectPreviousArtifactReadiness(
ctx context.Context,
cfg *config.Config,
store storage.ObjectStore,
requirements []artifacts.PreviousArtifactRequirement,
) previousArtifactReadiness {
out := previousArtifactReadiness{
Requirements: append([]artifacts.PreviousArtifactRequirement(nil), requirements...),
}
if len(requirements) == 0 {
return out
}
if strings.TrimSpace(cfg.Session.PreviousSessionID) == "" {
out.MissingID = true
return out
}
if store == nil {
out.Err = fmt.Errorf("previous-session artifacts cannot be checked because storage is unavailable")
return out
}
prefix := artifacts.S3SessionPrefix(cfg.Pipeline.Storage.S3.RootPrefix, cfg.Session.Campaign, cfg.Session.PreviousSessionID)
if _, err := artifacts.LoadCurrentState(ctx, store, prefix, artifacts.CurrentStateValidation{
ExpectedSessionID: strings.TrimSpace(cfg.Session.PreviousSessionID),
ExpectedCampaign: strings.TrimSpace(cfg.Session.Campaign),
ValidateRunID: true,
}); err != nil {
out.Err = fmt.Errorf("remote %v", err)
}
return out
}
func inspectRemoteCurrentState(ctx context.Context, cfg *config.Config, store storage.ObjectStore) remoteCurrentStateCheck {
if store == nil {
return remoteCurrentStateCheck{}
}
current, err := discoverRemoteCurrentStateFn(ctx, cfg, store)
if err != nil {
return remoteCurrentStateCheck{Err: err}
}
return remoteCurrentStateCheck{State: current}
}
func inspectEffectiveLocks(ctx context.Context, cfg *config.Config, store storage.ObjectStore) effectiveLocksCheck {
locks, err := loadEffectiveLocks(ctx, cfg, store)
if err != nil {
return effectiveLocksCheck{Err: err}
}
return effectiveLocksCheck{Locks: locks}
}
func resolveLocalInspectionAudioPaths(sessionDir string, inputs config.SessionInputsConfig) ([]string, error) {
if len(inputs.AudioFiles) > 0 {
out := make([]string, 0, len(inputs.AudioFiles))
seenBase := map[string]string{}
for _, item := range inputs.AudioFiles {
resolved, err := resolveInspectionPath(sessionDir, item)
if err != nil {
return nil, err
}
if !isInspectionFlacPath(resolved) {
return nil, fmt.Errorf("audio file %q must have .flac extension", resolved)
}
if err := requireInspectionFile(resolved, "audio file"); err != nil {
return nil, err
}
base := filepath.Base(resolved)
if prev, exists := seenBase[base]; exists && prev != resolved {
return nil, fmt.Errorf("duplicate audio basename %q from %q and %q", base, prev, resolved)
}
seenBase[base] = resolved
out = append(out, resolved)
}
sort.Strings(out)
return out, nil
}
audioDir, err := resolveInspectionPath(sessionDir, inputs.AudioDir)
if err != nil {
return nil, err
}
entries, err := os.ReadDir(audioDir)
if err != nil {
return nil, fmt.Errorf("read audio directory %q: %w", audioDir, err)
}
out := make([]string, 0, len(entries))
for _, entry := range entries {
if entry.IsDir() {
continue
}
full := filepath.Join(audioDir, entry.Name())
if !isInspectionFlacPath(full) {
continue
}
if err := requireInspectionFile(full, "audio file"); err != nil {
return nil, err
}
out = append(out, full)
}
if len(out) == 0 {
return nil, fmt.Errorf("no .flac files found in audio directory %q", audioDir)
}
sort.Strings(out)
return out, nil
}
func resolveInspectionPath(baseDir, inputPath string) (string, error) {
pathValue := strings.TrimSpace(inputPath)
if pathValue == "" {
return "", fmt.Errorf("path is required")
}
if filepath.IsAbs(pathValue) {
return filepath.Clean(pathValue), nil
}
return filepath.Clean(filepath.Join(baseDir, pathValue)), nil
}
func requireInspectionFile(path, label string) error {
info, err := os.Stat(path)
if err != nil {
if os.IsNotExist(err) {
return fmt.Errorf("%s %q does not exist", label, path)
}
return fmt.Errorf("stat %s %q: %w", label, path, err)
}
if info.IsDir() {
return fmt.Errorf("%s %q is a directory", label, path)
}
return nil
}
func isInspectionFlacPath(path string) bool {
return strings.EqualFold(filepath.Ext(strings.TrimSpace(path)), ".flac")
}

View File

@@ -48,13 +48,6 @@ func LocksList(ctx context.Context, args []string, out io.Writer) error {
// LocksAdd adds or updates one remote lock. // LocksAdd adds or updates one remote lock.
func LocksAdd(ctx context.Context, args []string, out io.Writer) error { func LocksAdd(ctx context.Context, args []string, out io.Writer) error {
var positionalSessionID string
var source string
if len(args) >= 2 && !isCLIFlagToken(args[0]) && !isCLIFlagToken(args[1]) {
positionalSessionID = strings.TrimSpace(args[0])
source = strings.TrimSpace(args[1])
args = append([]string(nil), args[2:]...)
}
fs := flag.NewFlagSet("locks add", flag.ContinueOnError) fs := flag.NewFlagSet("locks add", flag.ContinueOnError)
fs.SetOutput(io.Discard) fs.SetOutput(io.Discard)
var flags commonConfigFlags var flags commonConfigFlags
@@ -63,26 +56,8 @@ func LocksAdd(ctx context.Context, args []string, out io.Writer) error {
addCommonConfigFlags(fs, &flags) addCommonConfigFlags(fs, &flags)
fs.StringVar(&reason, "reason", "", "lock reason") fs.StringVar(&reason, "reason", "", "lock reason")
fs.BoolVar(&force, "force", false, "update existing remote lock") fs.BoolVar(&force, "force", false, "update existing remote lock")
if err := fs.Parse(args); err != nil { source, err := parseSessionIDAndOnePositionalArg("locks add", "source id", fs, args, &flags.sessionID)
return fmt.Errorf("locks add: invalid flags: %w", err) if err != nil {
}
if source == "" {
switch fs.NArg() {
case 2:
positionalSessionID = strings.TrimSpace(fs.Arg(0))
source = strings.TrimSpace(fs.Arg(1))
case 1:
if strings.TrimSpace(flags.sessionID) == "" {
return fmt.Errorf("locks add: expected session_id and source id")
}
source = strings.TrimSpace(fs.Arg(0))
default:
return fmt.Errorf("locks add: expected session_id and source id")
}
} else if fs.NArg() != 0 {
return fmt.Errorf("locks add: unexpected positional arguments")
}
if err := applyPositionalSessionID("locks add", positionalSessionID, &flags.sessionID); err != nil {
return err return err
} }
if strings.TrimSpace(flags.sessionID) == "" { if strings.TrimSpace(flags.sessionID) == "" {
@@ -116,37 +91,12 @@ func LocksAdd(ctx context.Context, args []string, out io.Writer) error {
// LocksRemove removes one remote lock. // LocksRemove removes one remote lock.
func LocksRemove(ctx context.Context, args []string, out io.Writer) error { func LocksRemove(ctx context.Context, args []string, out io.Writer) error {
var positionalSessionID string
var source string
if len(args) >= 2 && !isCLIFlagToken(args[0]) && !isCLIFlagToken(args[1]) {
positionalSessionID = strings.TrimSpace(args[0])
source = strings.TrimSpace(args[1])
args = append([]string(nil), args[2:]...)
}
fs := flag.NewFlagSet("locks remove", flag.ContinueOnError) fs := flag.NewFlagSet("locks remove", flag.ContinueOnError)
fs.SetOutput(io.Discard) fs.SetOutput(io.Discard)
var flags commonConfigFlags var flags commonConfigFlags
addCommonConfigFlags(fs, &flags) addCommonConfigFlags(fs, &flags)
if err := fs.Parse(args); err != nil { source, err := parseSessionIDAndOnePositionalArg("locks remove", "source id", fs, args, &flags.sessionID)
return fmt.Errorf("locks remove: invalid flags: %w", err) if err != nil {
}
if source == "" {
switch fs.NArg() {
case 2:
positionalSessionID = strings.TrimSpace(fs.Arg(0))
source = strings.TrimSpace(fs.Arg(1))
case 1:
if strings.TrimSpace(flags.sessionID) == "" {
return fmt.Errorf("locks remove: expected session_id and source id")
}
source = strings.TrimSpace(fs.Arg(0))
default:
return fmt.Errorf("locks remove: expected session_id and source id")
}
} else if fs.NArg() != 0 {
return fmt.Errorf("locks remove: unexpected positional arguments")
}
if err := applyPositionalSessionID("locks remove", positionalSessionID, &flags.sessionID); err != nil {
return err return err
} }
if strings.TrimSpace(flags.sessionID) == "" { if strings.TrimSpace(flags.sessionID) == "" {

View File

@@ -54,23 +54,26 @@ func SessionValidate(ctx context.Context, args []string, out io.Writer) error {
} }
requirements := artifacts.CollectPreviousArtifactRequirements(configuredScriptoriumArtifacts(cfg)) requirements := artifacts.CollectPreviousArtifactRequirements(configuredScriptoriumArtifacts(cfg))
if len(requirements) == 0 { previous := inspectPreviousArtifactReadiness(ctx, cfg, store, requirements)
if len(previous.Requirements) == 0 {
findings = append(findings, okFinding("previous", "no previous-session artifacts required")) findings = append(findings, okFinding("previous", "no previous-session artifacts required"))
} else if strings.TrimSpace(cfg.Session.PreviousSessionID) == "" { } else if previous.MissingID {
findings = append(findings, errorFinding("previous", "previous_session_id is required by configured previous-session artifacts")) findings = append(findings, errorFinding("previous", "previous_session_id is required by configured previous-session artifacts"))
} else if storeErr != nil { } else if previous.Err != nil {
findings = append(findings, errorFinding("previous", "previous-session artifacts cannot be checked because storage is unavailable")) findings = append(findings, errorFinding("previous", previous.Err.Error()))
} else { } else {
findings = append(findings, validatePreviousArtifactFindings(ctx, cfg, store, requirements)...) for _, req := range previous.Requirements {
findings = append(findings, okFinding("previous", fmt.Sprintf("%s required=%t", req.Name, req.Required)))
}
} }
locks, lockErr := loadEffectiveLocks(ctx, cfg, store) locks := inspectEffectiveLocks(ctx, cfg, store)
if lockErr != nil { if locks.Err != nil {
findings = append(findings, errorFinding("locks", lockErr.Error())) findings = append(findings, errorFinding("locks", locks.Err.Error()))
} else if len(locks.All) == 0 { } else if len(locks.Locks.All) == 0 {
findings = append(findings, okFinding("locks", "no effective publish locks")) findings = append(findings, okFinding("locks", "no effective publish locks"))
} else { } else {
for _, lock := range locks.All { for _, lock := range locks.Locks.All {
findings = append(findings, warnFinding("locks", fmt.Sprintf("%s locked: %s", lock.Source, strings.TrimSpace(lock.Reason)))) findings = append(findings, warnFinding("locks", fmt.Sprintf("%s locked: %s", lock.Source, strings.TrimSpace(lock.Reason))))
} }
} }

View File

@@ -5,8 +5,10 @@ import (
"flag" "flag"
"fmt" "fmt"
"io" "io"
"sort"
"strings" "strings"
"gitea.maximumdirect.net/eric/narratio/internal/adapters/storage"
"gitea.maximumdirect.net/eric/narratio/internal/artifacts" "gitea.maximumdirect.net/eric/narratio/internal/artifacts"
"gitea.maximumdirect.net/eric/narratio/internal/config" "gitea.maximumdirect.net/eric/narratio/internal/config"
) )
@@ -35,6 +37,8 @@ func Status(ctx context.Context, args []string, out io.Writer) error {
fmt.Fprintf(out, "Campaign: %s\n", cfg.Session.Campaign) fmt.Fprintf(out, "Campaign: %s\n", cfg.Session.Campaign)
fmt.Fprintf(out, "Workspace: %s\n", paths.Root) fmt.Fprintf(out, "Workspace: %s\n", paths.Root)
fmt.Fprintf(out, "Session config: %s\n", sessionSourceSummary(cfg)) fmt.Fprintf(out, "Session config: %s\n", sessionSourceSummary(cfg))
writeStatusStableInputs(out, inspectStableInputs(cfg))
writeStatusLocalAudio(out, inspectLocalAudioPresence(cfg))
if m, err := loadLocalManifest(ctx, paths.ManifestPath); err != nil { if m, err := loadLocalManifest(ctx, paths.ManifestPath); err != nil {
fmt.Fprintf(out, "Local manifest: error: %v\n", err) fmt.Fprintf(out, "Local manifest: error: %v\n", err)
@@ -49,21 +53,30 @@ func Status(ctx context.Context, args []string, out io.Writer) error {
if storeErr != nil { if storeErr != nil {
fmt.Fprintf(out, "Remote publish: unavailable: %v\n", storeErr) fmt.Fprintf(out, "Remote publish: unavailable: %v\n", storeErr)
} else if store != nil { } else if store != nil {
current, err := discoverRemoteCurrentStateFn(ctx, cfg, store) current := inspectRemoteCurrentState(ctx, cfg, store)
if err != nil { if current.Err != nil {
fmt.Fprintf(out, "Remote publish: missing or unavailable: %v\n", err) fmt.Fprintf(out, "Remote publish: missing or unavailable: %v\n", current.Err)
} else { } else {
fmt.Fprintf(out, "Remote publish: current run %s\n", current.RunID) fmt.Fprintf(out, "Remote publish: current run %s\n", current.State.RunID)
fmt.Fprintf(out, "Remote manifest: %s\n", current.CurrentManifestKey) fmt.Fprintf(out, "Remote manifest: %s\n", current.State.CurrentManifestKey)
} }
} }
writeStatusRemoteAudio(ctx, out, cfg, store, storeErr)
writeStatusPreviousArtifacts(out, inspectPreviousArtifactReadiness(
ctx,
cfg,
store,
artifacts.CollectPreviousArtifactRequirements(configuredScriptoriumArtifacts(cfg)),
))
locks, err := loadEffectiveLocks(ctx, cfg, store) lockChecks := inspectEffectiveLocks(ctx, cfg, store)
locks := lockChecks.Locks
lockErr := lockChecks.Err
if catalog, catalogErr := buildHelperArtifactCatalog(cfg); catalogErr != nil { if catalog, catalogErr := buildHelperArtifactCatalog(cfg); catalogErr != nil {
fmt.Fprintf(out, "Remote outputs: error: %v\n", catalogErr) fmt.Fprintf(out, "Remote outputs: error: %v\n", catalogErr)
} else if storeErr == nil { } else if storeErr == nil {
catalogLocks := locks catalogLocks := locks
if err != nil { if lockErr != nil {
catalogLocks = &effectiveLocks{ catalogLocks = &effectiveLocks{
Static: staticPublishLocks(cfg), Static: staticPublishLocks(cfg),
All: staticPublishLocks(cfg), All: staticPublishLocks(cfg),
@@ -76,8 +89,8 @@ func Status(ctx context.Context, args []string, out io.Writer) error {
fmt.Fprintln(out, "Remote outputs:") fmt.Fprintln(out, "Remote outputs:")
writeArtifactList(out, cfg, catalog, catalogLocks, publishedRemoteState) writeArtifactList(out, cfg, catalog, catalogLocks, publishedRemoteState)
} }
if err != nil { if lockErr != nil {
fmt.Fprintf(out, "Publish locks: error: %v\n", err) fmt.Fprintf(out, "Publish locks: error: %v\n", lockErr)
} else { } else {
writeLocks(out, cfg, locks) writeLocks(out, cfg, locks)
} }
@@ -86,3 +99,68 @@ func Status(ctx context.Context, args []string, out io.Writer) error {
fmt.Fprintf(out, "- narratio session restore %s --dry-run\n", cfg.Session.SessionID) fmt.Fprintf(out, "- narratio session restore %s --dry-run\n", cfg.Session.SessionID)
return nil return nil
} }
func writeStatusStableInputs(out io.Writer, checks []stableInputCheck) {
if len(checks) == 0 {
return
}
for _, check := range checks {
if check.Err != nil {
if strings.TrimSpace(check.Path) != "" {
fmt.Fprintf(out, "Stable input %s: unavailable: %v\n", check.Name, check.Err)
} else {
fmt.Fprintf(out, "Stable input %s: unavailable: %s\n", check.Name, check.Err.Error())
}
continue
}
fmt.Fprintf(out, "Stable input %s: %s\n", check.Name, check.Path)
}
}
func writeStatusLocalAudio(out io.Writer, check localAudioCheck) {
if !check.Checked {
return
}
if check.Err != nil {
fmt.Fprintf(out, "Local audio: unavailable: %v\n", check.Err)
return
}
fmt.Fprintf(out, "Local audio: %d file(s)\n", len(check.Paths))
}
func writeStatusRemoteAudio(ctx context.Context, out io.Writer, cfg *config.Config, store storage.ObjectStore, storeErr error) {
if cfg.Session.Inputs.AudioS3 == nil {
return
}
if storeErr != nil {
fmt.Fprintf(out, "Remote audio: unavailable: %v\n", storeErr)
return
}
check := inspectRemoteAudioPresence(ctx, cfg, store)
if check.Err != nil {
fmt.Fprintf(out, "Remote audio: unavailable: %v\n", check.Err)
return
}
fmt.Fprintf(out, "Remote audio: %d .flac object(s)\n", len(check.Keys))
}
func writeStatusPreviousArtifacts(out io.Writer, readiness previousArtifactReadiness) {
if len(readiness.Requirements) == 0 {
fmt.Fprintln(out, "Previous-session artifacts: not required")
return
}
if readiness.MissingID {
fmt.Fprintln(out, "Previous-session artifacts: unavailable: previous_session_id is required by configured previous-session artifacts")
return
}
if readiness.Err != nil {
fmt.Fprintf(out, "Previous-session artifacts: unavailable: %v\n", readiness.Err)
return
}
names := make([]string, 0, len(readiness.Requirements))
for _, req := range readiness.Requirements {
names = append(names, fmt.Sprintf("%s(required=%t)", req.Name, req.Required))
}
sort.Strings(names)
fmt.Fprintf(out, "Previous-session artifacts: ready: %s\n", strings.Join(names, ", "))
}

View File

@@ -11,6 +11,7 @@ import (
"gitea.maximumdirect.net/eric/narratio/internal/artifacts" "gitea.maximumdirect.net/eric/narratio/internal/artifacts"
"gitea.maximumdirect.net/eric/narratio/internal/audio" "gitea.maximumdirect.net/eric/narratio/internal/audio"
"gitea.maximumdirect.net/eric/narratio/internal/config" "gitea.maximumdirect.net/eric/narratio/internal/config"
"gitea.maximumdirect.net/eric/narratio/internal/fileops"
"gitea.maximumdirect.net/eric/narratio/internal/manifest" "gitea.maximumdirect.net/eric/narratio/internal/manifest"
) )
@@ -114,10 +115,7 @@ func executeRestoreDownloadAction(
} }
} }
if err := os.Chmod(tmpPath, 0o644); err != nil { if err := fileops.InstallDownloadedTempFile(tmpPath, safeLocalPath, 0o644); err != nil {
return fmt.Errorf("set file permissions: %w", err)
}
if err := os.Rename(tmpPath, safeLocalPath); err != nil {
return fmt.Errorf("install file atomically: %w", err) return fmt.Errorf("install file atomically: %w", err)
} }
removeTmp = false removeTmp = false

View File

@@ -2,6 +2,7 @@ package app
import ( import (
"context" "context"
"errors"
"fmt" "fmt"
"io" "io"
"os" "os"
@@ -13,6 +14,7 @@ import (
"gitea.maximumdirect.net/eric/narratio/internal/adapters/storage" "gitea.maximumdirect.net/eric/narratio/internal/adapters/storage"
"gitea.maximumdirect.net/eric/narratio/internal/artifacts" "gitea.maximumdirect.net/eric/narratio/internal/artifacts"
"gitea.maximumdirect.net/eric/narratio/internal/config" "gitea.maximumdirect.net/eric/narratio/internal/config"
"gitea.maximumdirect.net/eric/narratio/internal/pathsafe"
"gitea.maximumdirect.net/eric/narratio/internal/previouscache" "gitea.maximumdirect.net/eric/narratio/internal/previouscache"
) )
@@ -215,19 +217,17 @@ func joinWithinSessionRoot(sessionRoot, relative string) (string, error) {
if strings.TrimSpace(sessionRoot) == "" { if strings.TrimSpace(sessionRoot) == "" {
return "", fmt.Errorf("session root is required") return "", fmt.Errorf("session root is required")
} }
cleanRel := path.Clean(strings.TrimSpace(relative)) joined, err := pathsafe.JoinSlashRelativeUnderRoot(sessionRoot, filepath.ToSlash(strings.TrimSpace(relative)))
if cleanRel == "." || cleanRel == "" { if err != nil {
if errors.Is(err, pathsafe.ErrRelativePathRequired) {
return "", fmt.Errorf("relative path is required") return "", fmt.Errorf("relative path is required")
} }
if cleanRel == ".." || strings.HasPrefix(cleanRel, "../") || strings.HasPrefix(cleanRel, "/") { if errors.Is(err, pathsafe.ErrRelativePathEscape) || errors.Is(err, pathsafe.ErrRelativePathAbsolute) {
return "", fmt.Errorf("relative path escapes session root") return "", fmt.Errorf("relative path escapes session root")
} }
abs := filepath.Clean(filepath.Join(sessionRoot, filepath.FromSlash(cleanRel))) return "", fmt.Errorf("join relative path under session root: %w", err)
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 return joined, nil
} }
func buildPreviousCacheRestoreActions( func buildPreviousCacheRestoreActions(

View File

@@ -1,102 +0,0 @@
package app
import (
"context"
"flag"
"fmt"
"io"
"gitea.maximumdirect.net/eric/narratio/internal/artifacts"
"gitea.maximumdirect.net/eric/narratio/internal/config"
"gitea.maximumdirect.net/eric/narratio/internal/manifest"
)
// Resume continues execution from the first non-succeeded stage in the manifest.
func Resume(ctx context.Context, args []string, out io.Writer) error {
fs := flag.NewFlagSet("resume", flag.ContinueOnError)
fs.SetOutput(io.Discard)
var flags commonConfigFlags
var force bool
var selectedArtifacts artifactSelectionFlag
addCommonConfigFlags(fs, &flags)
fs.BoolVar(&force, "force", false, "force stage execution")
fs.Var(&selectedArtifacts, "artifacts", "configured artifact names to execute and publish (comma-separated or repeatable)")
if err := parseSessionAwareFlags("resume", fs, args, &flags.sessionID); err != nil {
return err
}
if flags.sessionID == "" {
return fmt.Errorf("resume: session_id is required")
}
cfg, err := loadCommandConfig(ctx, flags.pipelinePath, flags.campaignPath, flags.campaignFilePath, flags.sessionPath, flags.sessionOptions())
if err != nil {
return fmt.Errorf("resume: %w", err)
}
if err := config.Validate(cfg); err != nil {
return fmt.Errorf("resume: %w", err)
}
normalizedArtifacts, err := selectedArtifacts.Normalize()
if err != nil {
return fmt.Errorf("resume: invalid --artifacts: %w", err)
}
if err := validateSelectedArtifacts(cfg, normalizedArtifacts); err != nil {
return fmt.Errorf("resume: %w", err)
}
full := BuildFullPlan()
selected := full
if !force {
m, err := loadManifestIfPresent(ctx, cfg)
if err != nil {
return fmt.Errorf("resume: %w", err)
}
if m != nil {
start := firstNonSucceededIndex(full, m)
if start >= len(full) {
_, err := fmt.Fprintf(out, "narratio resume: session %s has no remaining stages\n", cfg.Session.SessionID)
return err
}
selected = full[start:]
}
}
summary, err := executeStagesFn(ctx, cfg, selected, RunOptions{
Force: force,
SelectedArtifacts: normalizedArtifacts,
})
if err != nil {
return fmt.Errorf("resume: %w", err)
}
_, err = fmt.Fprintf(
out,
"narratio resume: session %s; executed=%d skipped=%d; manifest=%s\n",
summary.SessionID,
len(summary.Executed),
len(summary.Skipped),
summary.ManifestPath,
)
return err
}
func loadManifestIfPresent(ctx context.Context, cfg *config.Config) (*manifest.Manifest, error) {
path := artifacts.SessionManifestPathForCampaign(
cfg.Pipeline.Workspace.Root,
cfg.Session.Campaign,
cfg.Session.SessionID,
)
exists, err := fileExists(path)
if err != nil {
return nil, fmt.Errorf("check manifest %q: %w", path, err)
}
if !exists {
return nil, nil
}
store := &manifest.LocalStore{}
m, err := store.Load(ctx, path)
if err != nil {
return nil, fmt.Errorf("load manifest %q: %w", path, err)
}
return m, nil
}

View File

@@ -1,8 +1,12 @@
package app package app
import ( import (
"context"
"fmt"
"time" "time"
"gitea.maximumdirect.net/eric/narratio/internal/artifacts"
"gitea.maximumdirect.net/eric/narratio/internal/config"
"gitea.maximumdirect.net/eric/narratio/internal/manifest" "gitea.maximumdirect.net/eric/narratio/internal/manifest"
"gitea.maximumdirect.net/eric/narratio/internal/stage" "gitea.maximumdirect.net/eric/narratio/internal/stage"
) )
@@ -43,13 +47,25 @@ func stageSucceeded(m *manifest.Manifest, name string) bool {
return sr != nil && sr.Status == manifest.StatusSucceeded return sr != nil && sr.Status == manifest.StatusSucceeded
} }
func firstNonSucceededIndex(stages []stage.Stage, m *manifest.Manifest) int { func loadManifestIfPresent(ctx context.Context, cfg *config.Config) (*manifest.Manifest, error) {
for i, s := range stages { path := artifacts.SessionManifestPathForCampaign(
if !stageSucceeded(m, s.Name()) { cfg.Pipeline.Workspace.Root,
return i cfg.Session.Campaign,
cfg.Session.SessionID,
)
exists, err := fileExists(path)
if err != nil {
return nil, fmt.Errorf("check manifest %q: %w", path, err)
} }
if !exists {
return nil, nil
} }
return len(stages) store := &manifest.LocalStore{}
m, err := store.Load(ctx, path)
if err != nil {
return nil, fmt.Errorf("load manifest %q: %w", path, err)
}
return m, nil
} }
func canonicalStageNames() []string { func canonicalStageNames() []string {

View File

@@ -8,18 +8,6 @@ import (
"gitea.maximumdirect.net/eric/narratio/internal/manifest" "gitea.maximumdirect.net/eric/narratio/internal/manifest"
) )
func TestFirstNonSucceededIndex(t *testing.T) {
stages := BuildFullPlan()
m := manifest.New("2026-05-03", time.Now().UTC())
m.MarkStageSucceeded("prepare", time.Now().UTC(), nil)
m.MarkStageSucceeded("transcribe", time.Now().UTC(), nil)
got := firstNonSucceededIndex(stages, m)
if got != 2 {
t.Fatalf("firstNonSucceededIndex() = %d, want 2", got)
}
}
func TestDecideStageActions(t *testing.T) { func TestDecideStageActions(t *testing.T) {
stages := BuildFullPlan()[:2] stages := BuildFullPlan()[:2]
m := manifest.New("2026-05-03", time.Now().UTC()) m := manifest.New("2026-05-03", time.Now().UTC())

View File

@@ -13,7 +13,7 @@ import (
"gitea.maximumdirect.net/eric/narratio/internal/manifest" "gitea.maximumdirect.net/eric/narratio/internal/manifest"
) )
func TestResumeStartsAfterCompletedStages(t *testing.T) { func TestRunContinuesAfterCompletedStages(t *testing.T) {
workspaceRoot := t.TempDir() workspaceRoot := t.TempDir()
pipelinePath, campaignPath, sessionPath := writeValidConfigFiles(t, workspaceRoot) pipelinePath, campaignPath, sessionPath := writeValidConfigFiles(t, workspaceRoot)
manifestPath := filepath.Join(workspaceRoot, "work", "sample-campaign", "2026-05-03", "manifest.json") manifestPath := filepath.Join(workspaceRoot, "work", "sample-campaign", "2026-05-03", "manifest.json")
@@ -32,12 +32,12 @@ func TestResumeStartsAfterCompletedStages(t *testing.T) {
mustWriteTestFile(t, filepath.Join(workRoot, "inputs", "glossary.yml"), "terms: []\n") mustWriteTestFile(t, filepath.Join(workRoot, "inputs", "glossary.yml"), "terms: []\n")
var out bytes.Buffer var out bytes.Buffer
err := Resume(context.Background(), []string{"2026-05-03", "--config", pipelinePath, "--campaign-file", campaignPath, "--session", sessionPath}, &out) err := Run(context.Background(), []string{"2026-05-03", "--config", pipelinePath, "--campaign-file", campaignPath, "--session", sessionPath}, &out)
if err != nil { if err != nil {
t.Fatalf("Resume() error = %v", err) t.Fatalf("Run() error = %v", err)
} }
if !strings.Contains(out.String(), "executed=7 skipped=0") { if !strings.Contains(out.String(), "executed=7 skipped=2") {
t.Fatalf("output = %q, want executed=7 skipped=0", out.String()) t.Fatalf("output = %q, want executed=7 skipped=2", out.String())
} }
loaded, err := store.Load(context.Background(), manifestPath) loaded, err := store.Load(context.Background(), manifestPath)
@@ -49,7 +49,7 @@ func TestResumeStartsAfterCompletedStages(t *testing.T) {
} }
} }
func TestResumeNoRemainingStages(t *testing.T) { func TestRunNoRemainingStagesRecordsSkippedStages(t *testing.T) {
workspaceRoot := t.TempDir() workspaceRoot := t.TempDir()
pipelinePath, campaignPath, sessionPath := writeValidConfigFiles(t, workspaceRoot) pipelinePath, campaignPath, sessionPath := writeValidConfigFiles(t, workspaceRoot)
manifestPath := filepath.Join(workspaceRoot, "work", "sample-campaign", "2026-05-03", "manifest.json") manifestPath := filepath.Join(workspaceRoot, "work", "sample-campaign", "2026-05-03", "manifest.json")
@@ -64,20 +64,20 @@ func TestResumeNoRemainingStages(t *testing.T) {
} }
var out bytes.Buffer var out bytes.Buffer
err := Resume(context.Background(), []string{"2026-05-03", "--config", pipelinePath, "--campaign-file", campaignPath, "--session", sessionPath}, &out) err := Run(context.Background(), []string{"2026-05-03", "--config", pipelinePath, "--campaign-file", campaignPath, "--session", sessionPath}, &out)
if err != nil { if err != nil {
t.Fatalf("Resume() error = %v", err) t.Fatalf("Run() error = %v", err)
} }
if !strings.Contains(out.String(), "has no remaining stages") { if !strings.Contains(out.String(), "executed=0 skipped=9") {
t.Fatalf("output = %q, want no remaining stages", out.String()) t.Fatalf("output = %q, want executed=0 skipped=9", out.String())
} }
} }
func TestResumeForceRerunsSucceeded(t *testing.T) { func TestRunForceRerunsSucceeded(t *testing.T) {
workspaceRoot := t.TempDir() workspaceRoot := t.TempDir()
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json") w.Header().Set("Content-Type", "application/json")
_, _ = w.Write([]byte(`{"source":"resume-force-test","segments":[{"speaker":"alice"}]}`)) _, _ = w.Write([]byte(`{"source":"run-force-test","segments":[{"speaker":"alice"}]}`))
})) }))
defer srv.Close() defer srv.Close()
pipelinePath, campaignPath, sessionPath := writeValidConfigFiles(t, workspaceRoot, srv.URL) pipelinePath, campaignPath, sessionPath := writeValidConfigFiles(t, workspaceRoot, srv.URL)
@@ -93,9 +93,9 @@ func TestResumeForceRerunsSucceeded(t *testing.T) {
} }
var out bytes.Buffer var out bytes.Buffer
err := Resume(context.Background(), []string{"2026-05-03", "--config", pipelinePath, "--campaign-file", campaignPath, "--session", sessionPath, "--force"}, &out) err := Run(context.Background(), []string{"2026-05-03", "--config", pipelinePath, "--campaign-file", campaignPath, "--session", sessionPath, "--force"}, &out)
if err != nil { if err != nil {
t.Fatalf("Resume() error = %v", err) t.Fatalf("Run() error = %v", err)
} }
if !strings.Contains(out.String(), "executed=9 skipped=0") { if !strings.Contains(out.String(), "executed=9 skipped=0") {
t.Fatalf("output = %q, want forced full rerun", out.String()) t.Fatalf("output = %q, want forced full rerun", out.String())
@@ -166,7 +166,7 @@ func TestRunStageSkipAndForce(t *testing.T) {
} }
} }
func TestRunStageForceMarksDownstreamStaleAndResumeContinuesFromStale(t *testing.T) { func TestRunStageForceMarksDownstreamStaleAndRunContinuesFromStale(t *testing.T) {
workspaceRoot := t.TempDir() workspaceRoot := t.TempDir()
pipelinePath, campaignPath, sessionPath := writeValidConfigFiles(t, workspaceRoot) pipelinePath, campaignPath, sessionPath := writeValidConfigFiles(t, workspaceRoot)
manifestPath := filepath.Join(workspaceRoot, "work", "sample-campaign", "2026-05-03", "manifest.json") manifestPath := filepath.Join(workspaceRoot, "work", "sample-campaign", "2026-05-03", "manifest.json")
@@ -203,12 +203,12 @@ func TestRunStageForceMarksDownstreamStaleAndResumeContinuesFromStale(t *testing
} }
out.Reset() out.Reset()
err = Resume(context.Background(), []string{"2026-05-03", "--config", pipelinePath, "--campaign-file", campaignPath, "--session", sessionPath}, &out) err = Run(context.Background(), []string{"2026-05-03", "--config", pipelinePath, "--campaign-file", campaignPath, "--session", sessionPath}, &out)
if err != nil { if err != nil {
t.Fatalf("Resume() error = %v", err) t.Fatalf("Run() error = %v", err)
} }
if !strings.Contains(out.String(), "executed=5 skipped=0") { if !strings.Contains(out.String(), "executed=5 skipped=4") {
t.Fatalf("output = %q, want resume to execute normalize..notify", out.String()) t.Fatalf("output = %q, want run to execute stale downstream stages", out.String())
} }
} }

View File

@@ -59,3 +59,37 @@ func parseSessionAwareFlags(command string, fs *flag.FlagSet, args []string, ses
} }
return resolveParsedSessionID(command, positionalSessionID, fs, sessionID) return resolveParsedSessionID(command, positionalSessionID, fs, sessionID)
} }
func parseSessionIDAndOnePositionalArg(command, argName string, fs *flag.FlagSet, args []string, sessionID *string) (string, error) {
var positionalSessionID string
value := ""
if len(args) >= 2 && !isCLIFlagToken(args[0]) && !isCLIFlagToken(args[1]) {
positionalSessionID = strings.TrimSpace(args[0])
value = strings.TrimSpace(args[1])
args = append([]string(nil), args[2:]...)
}
if err := fs.Parse(args); err != nil {
return "", fmt.Errorf("%s: invalid flags: %w", command, err)
}
if value == "" {
switch fs.NArg() {
case 2:
positionalSessionID = strings.TrimSpace(fs.Arg(0))
value = strings.TrimSpace(fs.Arg(1))
case 1:
if strings.TrimSpace(*sessionID) == "" {
return "", fmt.Errorf("%s: expected session_id and %s", command, argName)
}
value = strings.TrimSpace(fs.Arg(0))
default:
return "", fmt.Errorf("%s: expected session_id and %s", command, argName)
}
} else if fs.NArg() != 0 {
return "", fmt.Errorf("%s: unexpected positional arguments", command)
}
if err := applyPositionalSessionID(command, positionalSessionID, sessionID); err != nil {
return "", err
}
return value, nil
}

View File

@@ -162,8 +162,8 @@ func TestExecuteWorkflowCommandsAcceptPositionalSessionID(t *testing.T) {
wantForce bool wantForce bool
}{ }{
{ {
name: "resume", name: "run",
args: []string{"resume", "2026-05-03", "--config", pipelinePath, "--campaign-file", campaignPath, "--session", sessionPath}, args: []string{"run", "2026-05-03", "--config", pipelinePath, "--campaign-file", campaignPath, "--session", sessionPath},
wantStage: "prepare", wantStage: "prepare",
wantForce: false, wantForce: false,
}, },

View File

@@ -1,6 +1,7 @@
package artifactpolicy package artifactpolicy
import ( import (
"errors"
"fmt" "fmt"
"regexp" "regexp"
"strings" "strings"
@@ -19,6 +20,11 @@ const (
var configuredSourceRE = regexp.MustCompile(`^narratio\.artifact\.([a-z][a-z0-9_]*)$`) var configuredSourceRE = regexp.MustCompile(`^narratio\.artifact\.([a-z][a-z0-9_]*)$`)
var previousSourceRE = regexp.MustCompile(`^narratio\.previous_session\.artifact\.([a-z][a-z0-9_]*)$`) var previousSourceRE = regexp.MustCompile(`^narratio\.previous_session\.artifact\.([a-z][a-z0-9_]*)$`)
var (
ErrUnsupportedScriptoriumInputSource = errors.New("unsupported scriptorium input source")
ErrInvalidPreviousSessionSource = errors.New("invalid previous-session source format")
)
type SourceKind string type SourceKind string
const ( const (
@@ -34,6 +40,28 @@ type Source struct {
ConfiguredKey string ConfiguredKey string
} }
// ScriptoriumInputSourceDescriptor describes one validated Scriptorium input source.
type ScriptoriumInputSourceDescriptor struct {
Source Source
PreviousSession *PreviousSessionSourceDescriptor
}
// PreviousSessionSourceDescriptor describes one canonical previous-session input source.
type PreviousSessionSourceDescriptor struct {
SourceID string
ConfiguredKey string
ConfiguredSourceID string
}
// UnknownConfiguredArtifactError reports a source that references an undefined configured artifact key.
type UnknownConfiguredArtifactError struct {
ConfiguredKey string
}
func (e *UnknownConfiguredArtifactError) Error() string {
return fmt.Sprintf("references unknown artifact %q", e.ConfiguredKey)
}
// ConfiguredSourceID converts a configured artifact key into source id form. // ConfiguredSourceID converts a configured artifact key into source id form.
func ConfiguredSourceID(key string) string { func ConfiguredSourceID(key string) string {
return configuredSourcePrefix + strings.TrimSpace(key) return configuredSourcePrefix + strings.TrimSpace(key)
@@ -83,6 +111,70 @@ func ClassifySource(source string) (Source, error) {
return Source{}, fmt.Errorf("unsupported artifact source %q", source) return Source{}, fmt.Errorf("unsupported artifact source %q", source)
} }
// DescribeScriptoriumInputSource classifies one input source and returns descriptor
// metadata used by config validation, analyze input resolution, and previous-cache planning.
func DescribeScriptoriumInputSource(source string) (ScriptoriumInputSourceDescriptor, error) {
trimmed := strings.TrimSpace(source)
if trimmed == "" {
return ScriptoriumInputSourceDescriptor{}, ErrUnsupportedScriptoriumInputSource
}
if strings.HasPrefix(trimmed, "narratio.previous_session.artifact") {
descriptor, err := DescribePreviousSessionSource(trimmed)
if err != nil {
return ScriptoriumInputSourceDescriptor{}, err
}
return ScriptoriumInputSourceDescriptor{
Source: Source{
ID: descriptor.SourceID,
Kind: SourceKindPreviousArtifact,
ConfiguredKey: descriptor.ConfiguredKey,
},
PreviousSession: &descriptor,
}, nil
}
classified, err := ClassifySource(trimmed)
if err != nil {
return ScriptoriumInputSourceDescriptor{}, ErrUnsupportedScriptoriumInputSource
}
return ScriptoriumInputSourceDescriptor{Source: classified}, nil
}
// DescribePreviousSessionSource validates a canonical previous-session source id
// and returns both previous and configured-source vocabulary descriptors.
func DescribePreviousSessionSource(source string) (PreviousSessionSourceDescriptor, error) {
configuredKey, ok := ParsePreviousSessionSource(source)
if !ok {
return PreviousSessionSourceDescriptor{}, ErrInvalidPreviousSessionSource
}
return PreviousSessionSourceDescriptor{
SourceID: PreviousSessionSourceID(configuredKey),
ConfiguredKey: configuredKey,
ConfiguredSourceID: ConfiguredSourceID(configuredKey),
}, nil
}
// PreviousSessionSourceDescriptorForConfiguredKey derives a previous-session source descriptor
// from a configured artifact key.
func PreviousSessionSourceDescriptorForConfiguredKey(configuredKey string) (PreviousSessionSourceDescriptor, error) {
return DescribePreviousSessionSource(PreviousSessionSourceID(configuredKey))
}
// ValidateInputConfiguredReference checks that configured/previous-session sources
// reference configured artifacts known to the current Scriptorium config.
func ValidateInputConfiguredReference(
descriptor ScriptoriumInputSourceDescriptor,
configured map[string]struct{},
) error {
switch descriptor.Source.Kind {
case SourceKindConfiguredArtifact, SourceKindPreviousArtifact:
if _, ok := configured[descriptor.Source.ConfiguredKey]; !ok {
return &UnknownConfiguredArtifactError{ConfiguredKey: descriptor.Source.ConfiguredKey}
}
}
return nil
}
// ValidatePublishSource validates that a source is publish-compatible and references a known configured artifact. // ValidatePublishSource validates that a source is publish-compatible and references a known configured artifact.
func ValidatePublishSource(source string, configured map[string]string) (Source, error) { func ValidatePublishSource(source string, configured map[string]string) (Source, error) {
classified, err := ClassifySource(source) classified, err := ClassifySource(source)

View File

@@ -1,6 +1,7 @@
package artifactpolicy package artifactpolicy
import ( import (
"errors"
"strings" "strings"
"testing" "testing"
) )
@@ -90,3 +91,101 @@ func TestResolvePublishedDestinationRejectsTraversal(t *testing.T) {
t.Fatal("ResolvePublishedDestination() error = nil, want traversal rejection") t.Fatal("ResolvePublishedDestination() error = nil, want traversal rejection")
} }
} }
func TestDescribeScriptoriumInputSource(t *testing.T) {
tests := []struct {
name string
source string
wantKind SourceKind
wantKey string
wantPrev bool
wantErr error
wantErrLike string
}{
{name: "built in", source: "narratio.transcript.final_trimmed", wantKind: SourceKindBuiltIn},
{name: "configured", source: "narratio.artifact.session_recap", wantKind: SourceKindConfiguredArtifact, wantKey: "session_recap"},
{name: "previous", source: "narratio.previous_session.artifact.session_recap", wantKind: SourceKindPreviousArtifact, wantKey: "session_recap", wantPrev: true},
{name: "invalid previous", source: "narratio.previous_session.artifact.", wantErr: ErrInvalidPreviousSessionSource},
{name: "unsupported", source: "narratio.unknown", wantErr: ErrUnsupportedScriptoriumInputSource},
{name: "empty", source: " ", wantErr: ErrUnsupportedScriptoriumInputSource},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
got, err := DescribeScriptoriumInputSource(tt.source)
if tt.wantErr != nil {
if !errors.Is(err, tt.wantErr) {
t.Fatalf("DescribeScriptoriumInputSource() error = %v, want %v", err, tt.wantErr)
}
return
}
if tt.wantErrLike != "" {
if err == nil || !strings.Contains(err.Error(), tt.wantErrLike) {
t.Fatalf("DescribeScriptoriumInputSource() error = %v, want like %q", err, tt.wantErrLike)
}
return
}
if err != nil {
t.Fatalf("DescribeScriptoriumInputSource() error = %v", err)
}
if got.Source.Kind != tt.wantKind {
t.Fatalf("DescribeScriptoriumInputSource().Source.Kind = %q, want %q", got.Source.Kind, tt.wantKind)
}
if got.Source.ConfiguredKey != tt.wantKey {
t.Fatalf("DescribeScriptoriumInputSource().Source.ConfiguredKey = %q, want %q", got.Source.ConfiguredKey, tt.wantKey)
}
if tt.wantPrev && got.PreviousSession == nil {
t.Fatal("DescribeScriptoriumInputSource().PreviousSession = nil, want descriptor")
}
if !tt.wantPrev && got.PreviousSession != nil {
t.Fatalf("DescribeScriptoriumInputSource().PreviousSession = %#v, want nil", got.PreviousSession)
}
})
}
}
func TestValidateInputConfiguredReference(t *testing.T) {
configured := map[string]struct{}{"session_recap": {}}
desc, err := DescribeScriptoriumInputSource("narratio.artifact.session_recap")
if err != nil {
t.Fatalf("DescribeScriptoriumInputSource(configured) error = %v", err)
}
if err := ValidateInputConfiguredReference(desc, configured); err != nil {
t.Fatalf("ValidateInputConfiguredReference(configured) error = %v", err)
}
prevDesc, err := DescribeScriptoriumInputSource("narratio.previous_session.artifact.session_recap")
if err != nil {
t.Fatalf("DescribeScriptoriumInputSource(previous) error = %v", err)
}
if err := ValidateInputConfiguredReference(prevDesc, configured); err != nil {
t.Fatalf("ValidateInputConfiguredReference(previous) error = %v", err)
}
missingDesc, err := DescribeScriptoriumInputSource("narratio.artifact.quest_log")
if err != nil {
t.Fatalf("DescribeScriptoriumInputSource(missing configured) error = %v", err)
}
err = ValidateInputConfiguredReference(missingDesc, configured)
var unknown *UnknownConfiguredArtifactError
if !errors.As(err, &unknown) || unknown.ConfiguredKey != "quest_log" {
t.Fatalf("ValidateInputConfiguredReference(missing configured) error = %v, want UnknownConfiguredArtifactError(quest_log)", err)
}
}
func TestPreviousSessionSourceDescriptorForConfiguredKey(t *testing.T) {
got, err := PreviousSessionSourceDescriptorForConfiguredKey("session_recap")
if err != nil {
t.Fatalf("PreviousSessionSourceDescriptorForConfiguredKey() error = %v", err)
}
if got.SourceID != "narratio.previous_session.artifact.session_recap" {
t.Fatalf("SourceID = %q, want narratio.previous_session.artifact.session_recap", got.SourceID)
}
if got.ConfiguredSourceID != "narratio.artifact.session_recap" {
t.Fatalf("ConfiguredSourceID = %q, want narratio.artifact.session_recap", got.ConfiguredSourceID)
}
if got.ConfiguredKey != "session_recap" {
t.Fatalf("ConfiguredKey = %q, want session_recap", got.ConfiguredKey)
}
}

View File

@@ -9,6 +9,9 @@ import (
"strconv" "strconv"
"strings" "strings"
"time" "time"
"gitea.maximumdirect.net/eric/narratio/internal/fileops"
"gitea.maximumdirect.net/eric/narratio/internal/pathsafe"
) )
// ErrLockConflict is returned when a session lock already exists. // ErrLockConflict is returned when a session lock already exists.
@@ -101,7 +104,7 @@ func (s *LocalStore) copyInputWithPaths(paths SessionPaths, sessionID, srcPath,
return Ref{}, fmt.Errorf("copy input: %w", err) return Ref{}, fmt.Errorf("copy input: %w", err)
} }
if err := copyFileAtomic(srcPath, destAbs, 0o644); err != nil { if err := fileops.CopyFileAtomic(srcPath, destAbs, 0o644); err != nil {
return Ref{}, fmt.Errorf("copy input %q -> %q: %w", srcPath, destAbs, err) return Ref{}, fmt.Errorf("copy input %q -> %q: %w", srcPath, destAbs, err)
} }
@@ -145,45 +148,9 @@ func (s *LocalStore) WriteFileAtomic(path string, data []byte, perm os.FileMode)
if strings.TrimSpace(path) == "" { if strings.TrimSpace(path) == "" {
return fmt.Errorf("write file atomic: path is required") return fmt.Errorf("write file atomic: path is required")
} }
if err := fileops.WriteFileAtomic(path, data, perm); err != nil {
dir := filepath.Dir(path) return fmt.Errorf("write file atomic: %w", err)
if err := os.MkdirAll(dir, 0o755); err != nil {
return fmt.Errorf("write file atomic: create parent dir %q: %w", dir, err)
} }
base := filepath.Base(path)
tmp, err := os.CreateTemp(dir, "."+base+".tmp-*")
if err != nil {
return fmt.Errorf("write file atomic: create temp file: %w", err)
}
tmpName := tmp.Name()
removeTmp := true
defer func() {
if removeTmp {
_ = os.Remove(tmpName)
}
}()
if _, err := tmp.Write(data); err != nil {
_ = tmp.Close()
return fmt.Errorf("write file atomic: write temp file: %w", err)
}
if err := tmp.Sync(); err != nil {
_ = tmp.Close()
return fmt.Errorf("write file atomic: sync temp file: %w", err)
}
if err := tmp.Close(); err != nil {
return fmt.Errorf("write file atomic: close temp file: %w", err)
}
if err := os.Chmod(tmpName, perm); err != nil {
return fmt.Errorf("write file atomic: chmod temp file: %w", err)
}
if err := os.Rename(tmpName, path); err != nil {
return fmt.Errorf("write file atomic: rename temp file: %w", err)
}
removeTmp = false
return nil return nil
} }
@@ -256,62 +223,21 @@ func (s *LocalStore) ReleaseSessionLock(lock *LockHandle) error {
} }
func resolveInRoot(root, relative string) (string, error) { func resolveInRoot(root, relative string) (string, error) {
rel := filepath.Clean(relative) joined, err := pathsafe.JoinSlashRelativeUnderRoot(root, filepath.ToSlash(relative))
if rel == "." || rel == "" { if err != nil {
switch {
case errors.Is(err, pathsafe.ErrRelativePathRequired):
return "", fmt.Errorf("relative destination path is required")
case errors.Is(err, pathsafe.ErrRelativePathAbsolute):
return "", fmt.Errorf("relative destination must not be absolute: %q", relative)
case errors.Is(err, pathsafe.ErrRelativePathEscape):
return "", fmt.Errorf("relative destination escapes root: %q", relative)
default:
return "", fmt.Errorf("resolve destination in root: %w", err)
}
}
if strings.TrimSpace(joined) == "" {
return "", fmt.Errorf("relative destination path is required") return "", fmt.Errorf("relative destination path is required")
} }
if filepath.IsAbs(rel) { return joined, nil
return "", fmt.Errorf("relative destination must not be absolute: %q", relative)
}
if rel == ".." || strings.HasPrefix(rel, ".."+string(filepath.Separator)) {
return "", fmt.Errorf("relative destination escapes root: %q", relative)
}
return filepath.Join(root, rel), nil
}
func copyFileAtomic(srcPath, dstPath string, perm os.FileMode) error {
src, err := os.Open(srcPath)
if err != nil {
return err
}
defer src.Close()
if err := os.MkdirAll(filepath.Dir(dstPath), 0o755); err != nil {
return err
}
dir := filepath.Dir(dstPath)
base := filepath.Base(dstPath)
tmp, err := os.CreateTemp(dir, "."+base+".tmp-*")
if err != nil {
return err
}
tmpName := tmp.Name()
removeTmp := true
defer func() {
if removeTmp {
_ = os.Remove(tmpName)
}
}()
if _, err := io.Copy(tmp, src); err != nil {
_ = tmp.Close()
return err
}
if err := tmp.Sync(); err != nil {
_ = tmp.Close()
return err
}
if err := tmp.Close(); err != nil {
return err
}
if err := os.Chmod(tmpName, perm); err != nil {
return err
}
if err := os.Rename(tmpName, dstPath); err != nil {
return err
}
removeTmp = false
return nil
} }

View File

@@ -5,6 +5,7 @@ import (
"sort" "sort"
"strings" "strings"
"gitea.maximumdirect.net/eric/narratio/internal/artifactpolicy"
"gitea.maximumdirect.net/eric/narratio/internal/config" "gitea.maximumdirect.net/eric/narratio/internal/config"
) )
@@ -36,10 +37,11 @@ func CollectPreviousArtifactRequirements(
inputNames := sortedScriptoriumInputKeys(artifactCfg.Inputs) inputNames := sortedScriptoriumInputKeys(artifactCfg.Inputs)
for _, inputName := range inputNames { for _, inputName := range inputNames {
inputCfg := artifactCfg.Inputs[inputName] inputCfg := artifactCfg.Inputs[inputName]
previousName, ok := PreviousSessionArtifactName(inputCfg.Source) descriptor, err := artifactpolicy.DescribeScriptoriumInputSource(inputCfg.Source)
if !ok { if err != nil || descriptor.PreviousSession == nil {
continue continue
} }
previousName := descriptor.PreviousSession.ConfiguredKey
location := fmt.Sprintf( location := fmt.Sprintf(
"pipeline.scriptorium.artifacts.%s.inputs.%s.source", "pipeline.scriptorium.artifacts.%s.inputs.%s.source",

View File

@@ -2,16 +2,14 @@ package audio
import ( import (
"context" "context"
"crypto/sha256"
"encoding/hex"
"fmt" "fmt"
"io"
"os" "os"
"path/filepath" "path/filepath"
"strings" "strings"
"gitea.maximumdirect.net/eric/narratio/internal/adapters/storage" "gitea.maximumdirect.net/eric/narratio/internal/adapters/storage"
"gitea.maximumdirect.net/eric/narratio/internal/artifacts" "gitea.maximumdirect.net/eric/narratio/internal/artifacts"
"gitea.maximumdirect.net/eric/narratio/internal/fileops"
) )
// S3MaterializeRequest describes one S3-backed audio materialization. // S3MaterializeRequest describes one S3-backed audio materialization.
@@ -61,7 +59,7 @@ func MaterializeS3Audio(ctx context.Context, req S3MaterializeRequest) (S3Materi
if ok, err := validCachedAudio(cachePath, req.Object.Size); err != nil { if ok, err := validCachedAudio(cachePath, req.Object.Size); err != nil {
return S3MaterializeResult{}, err return S3MaterializeResult{}, err
} else if ok { } else if ok {
checksum, err := copyFileAtomicWithChecksum(cachePath, req.DestPath, 0o644) checksum, err := fileops.CopyFileAtomicWithChecksum(cachePath, req.DestPath, 0o644)
if err != nil { if err != nil {
return S3MaterializeResult{}, fmt.Errorf("materialize cached audio %q: %w", cachePath, err) return S3MaterializeResult{}, fmt.Errorf("materialize cached audio %q: %w", cachePath, err)
} }
@@ -82,7 +80,7 @@ func MaterializeS3Audio(ctx context.Context, req S3MaterializeRequest) (S3Materi
return S3MaterializeResult{}, fmt.Errorf("validate downloaded audio %q: %w", spoolPath, err) return S3MaterializeResult{}, fmt.Errorf("validate downloaded audio %q: %w", spoolPath, err)
} }
checksum, err := copyFileAtomicWithChecksum(spoolPath, req.DestPath, 0o644) checksum, err := fileops.CopyFileAtomicWithChecksum(spoolPath, req.DestPath, 0o644)
if err != nil { if err != nil {
return S3MaterializeResult{}, fmt.Errorf("materialize downloaded audio %q: %w", filepath.Base(req.DestPath), err) return S3MaterializeResult{}, fmt.Errorf("materialize downloaded audio %q: %w", filepath.Base(req.DestPath), err)
} }
@@ -91,7 +89,7 @@ func MaterializeS3Audio(ctx context.Context, req S3MaterializeRequest) (S3Materi
result.Downloaded = true result.Downloaded = true
if result.CachePath != "" { if result.CachePath != "" {
if _, err := copyFileAtomicWithChecksum(spoolPath, result.CachePath, 0o644); err != nil { if _, err := fileops.CopyFileAtomicWithChecksum(spoolPath, result.CachePath, 0o644); err != nil {
return S3MaterializeResult{}, fmt.Errorf("populate audio cache %q: %w", result.CachePath, err) return S3MaterializeResult{}, fmt.Errorf("populate audio cache %q: %w", result.CachePath, err)
} }
} }
@@ -164,60 +162,9 @@ func downloadObjectAtomic(ctx context.Context, store storage.ObjectStore, key, d
if err := store.Download(ctx, key, tmpPath); err != nil { if err := store.Download(ctx, key, tmpPath); err != nil {
return err return err
} }
if err := os.Chmod(tmpPath, 0o644); err != nil { if err := fileops.InstallDownloadedTempFile(tmpPath, destPath, 0o644); err != nil {
return fmt.Errorf("set temp file permissions: %w", err) return err
}
if err := os.Rename(tmpPath, destPath); err != nil {
return fmt.Errorf("install downloaded file: %w", err)
} }
removeTmp = false removeTmp = false
return nil return nil
} }
func copyFileAtomicWithChecksum(src, dst string, perm os.FileMode) (string, error) {
if strings.TrimSpace(src) == "" || strings.TrimSpace(dst) == "" {
return "", fmt.Errorf("source and destination paths are required")
}
in, err := os.Open(src)
if err != nil {
return "", err
}
defer func() { _ = in.Close() }()
if err := os.MkdirAll(filepath.Dir(dst), 0o755); err != nil {
return "", fmt.Errorf("create destination directory: %w", err)
}
base := filepath.Base(dst)
tmp, err := os.CreateTemp(filepath.Dir(dst), "."+base+".tmp-*")
if err != nil {
return "", fmt.Errorf("create temp file: %w", err)
}
tmpPath := tmp.Name()
removeTmp := true
defer func() {
if removeTmp {
_ = os.Remove(tmpPath)
}
}()
digest := sha256.New()
if _, err := io.Copy(io.MultiWriter(tmp, digest), in); err != nil {
_ = tmp.Close()
return "", fmt.Errorf("copy file: %w", err)
}
if err := tmp.Sync(); err != nil {
_ = tmp.Close()
return "", fmt.Errorf("sync temp file: %w", err)
}
if err := tmp.Close(); err != nil {
return "", fmt.Errorf("close temp file: %w", err)
}
if err := os.Chmod(tmpPath, perm); err != nil {
return "", fmt.Errorf("chmod temp file: %w", err)
}
if err := os.Rename(tmpPath, dst); err != nil {
return "", fmt.Errorf("install temp file: %w", err)
}
removeTmp = false
return hex.EncodeToString(digest.Sum(nil)), nil
}

View File

@@ -9,7 +9,6 @@ import (
"strings" "strings"
"time" "time"
"gitea.maximumdirect.net/eric/narratio/internal/artifactmodel"
"gitea.maximumdirect.net/eric/narratio/internal/artifactpolicy" "gitea.maximumdirect.net/eric/narratio/internal/artifactpolicy"
"gitea.maximumdirect.net/eric/narratio/internal/pathsafe" "gitea.maximumdirect.net/eric/narratio/internal/pathsafe"
) )
@@ -639,13 +638,9 @@ var scriptoriumArtifactKeyRE = regexp.MustCompile(`^[a-z][a-z0-9_]*$`)
func validateScriptoriumInputSource(artifactName, inputName, source string, configuredArtifacts map[string]struct{}) (string, error) { func validateScriptoriumInputSource(artifactName, inputName, source string, configuredArtifacts map[string]struct{}) (string, error) {
trimmedSource := strings.TrimSpace(source) trimmedSource := strings.TrimSpace(source)
if isStaticSupportedScriptoriumInputSource(trimmedSource) { descriptor, err := artifactpolicy.DescribeScriptoriumInputSource(trimmedSource)
return "", nil if err != nil {
} if errors.Is(err, artifactpolicy.ErrInvalidPreviousSessionSource) {
if strings.HasPrefix(trimmedSource, "narratio.previous_session.artifact") {
referenced, ok := artifactpolicy.ParsePreviousSessionSource(trimmedSource)
if !ok {
return "", fmt.Errorf( return "", fmt.Errorf(
"pipeline.scriptorium.artifacts.%s.inputs.%s.source %q must reference configured artifact key matching ^[a-z][a-z0-9_]*$", "pipeline.scriptorium.artifacts.%s.inputs.%s.source %q must reference configured artifact key matching ^[a-z][a-z0-9_]*$",
artifactName, artifactName,
@@ -653,20 +648,6 @@ func validateScriptoriumInputSource(artifactName, inputName, source string, conf
source, source,
) )
} }
if _, ok := configuredArtifacts[referenced]; !ok {
return "", fmt.Errorf(
"pipeline.scriptorium.artifacts.%s.inputs.%s.source %q references unknown artifact %q",
artifactName,
inputName,
source,
referenced,
)
}
return "", nil
}
referenced, ok := artifactpolicy.ParseConfiguredSource(trimmedSource)
if !ok {
return "", fmt.Errorf( return "", fmt.Errorf(
"pipeline.scriptorium.artifacts.%s.inputs.%s.source %q is unsupported", "pipeline.scriptorium.artifacts.%s.inputs.%s.source %q is unsupported",
artifactName, artifactName,
@@ -674,28 +655,28 @@ func validateScriptoriumInputSource(artifactName, inputName, source string, conf
source, source,
) )
} }
if _, ok := configuredArtifacts[referenced]; !ok { if err := artifactpolicy.ValidateInputConfiguredReference(descriptor, configuredArtifacts); err != nil {
var unknownConfigured *artifactpolicy.UnknownConfiguredArtifactError
if errors.As(err, &unknownConfigured) {
return "", fmt.Errorf( return "", fmt.Errorf(
"pipeline.scriptorium.artifacts.%s.inputs.%s.source %q references unknown artifact %q", "pipeline.scriptorium.artifacts.%s.inputs.%s.source %q references unknown artifact %q",
artifactName, artifactName,
inputName, inputName,
source, source,
referenced, unknownConfigured.ConfiguredKey,
) )
} }
return referenced, nil return "", fmt.Errorf(
"pipeline.scriptorium.artifacts.%s.inputs.%s.source %q is unsupported",
artifactName,
inputName,
source,
)
} }
if descriptor.Source.Kind == artifactpolicy.SourceKindConfiguredArtifact {
func isStaticSupportedScriptoriumInputSource(source string) bool { return descriptor.Source.ConfiguredKey, nil
if _, ok := artifactmodel.LookupRuntimeTranscriptArtifact(source); ok {
return true
}
switch source {
case "narratio.bounds.session":
return true
default:
return false
} }
return "", nil
} }
func validateEnvVarNameField(fieldName, value string) error { func validateEnvVarNameField(fieldName, value string) error {

128
internal/fileops/fileops.go Normal file
View File

@@ -0,0 +1,128 @@
package fileops
import (
"crypto/sha256"
"encoding/hex"
"fmt"
"io"
"os"
"path/filepath"
"strings"
)
// WriteFileAtomic writes data to dst atomically via temp file + rename.
func WriteFileAtomic(dst string, data []byte, perm os.FileMode) error {
if strings.TrimSpace(dst) == "" {
return fmt.Errorf("destination path is required")
}
if err := os.MkdirAll(filepath.Dir(dst), 0o755); err != nil {
return fmt.Errorf("create destination directory: %w", err)
}
base := filepath.Base(dst)
tmp, err := os.CreateTemp(filepath.Dir(dst), "."+base+".tmp-*")
if err != nil {
return fmt.Errorf("create temp file: %w", err)
}
tmpPath := tmp.Name()
removeTmp := true
defer func() {
if removeTmp {
_ = os.Remove(tmpPath)
}
}()
if _, err := tmp.Write(data); err != nil {
_ = tmp.Close()
return fmt.Errorf("write temp file: %w", err)
}
if err := tmp.Sync(); err != nil {
_ = tmp.Close()
return fmt.Errorf("sync temp file: %w", err)
}
if err := tmp.Close(); err != nil {
return fmt.Errorf("close temp file: %w", err)
}
if err := os.Chmod(tmpPath, perm); err != nil {
return fmt.Errorf("set temp file permissions: %w", err)
}
if err := os.Rename(tmpPath, dst); err != nil {
return fmt.Errorf("install temp file: %w", err)
}
removeTmp = false
return nil
}
// CopyFileAtomic copies src to dst atomically via temp file + rename.
func CopyFileAtomic(src, dst string, perm os.FileMode) error {
_, err := CopyFileAtomicWithChecksum(src, dst, perm)
return err
}
// CopyFileAtomicWithChecksum copies src to dst atomically and returns the SHA-256 checksum.
func CopyFileAtomicWithChecksum(src, dst string, perm os.FileMode) (string, error) {
if strings.TrimSpace(src) == "" || strings.TrimSpace(dst) == "" {
return "", fmt.Errorf("source and destination paths are required")
}
in, err := os.Open(src)
if err != nil {
return "", fmt.Errorf("open source file: %w", err)
}
defer func() { _ = in.Close() }()
if err := os.MkdirAll(filepath.Dir(dst), 0o755); err != nil {
return "", fmt.Errorf("create destination directory: %w", err)
}
base := filepath.Base(dst)
tmp, err := os.CreateTemp(filepath.Dir(dst), "."+base+".tmp-*")
if err != nil {
return "", fmt.Errorf("create temp file: %w", err)
}
tmpPath := tmp.Name()
removeTmp := true
defer func() {
if removeTmp {
_ = os.Remove(tmpPath)
}
}()
digest := sha256.New()
if _, err := io.Copy(io.MultiWriter(tmp, digest), in); err != nil {
_ = tmp.Close()
return "", fmt.Errorf("copy file: %w", err)
}
if err := tmp.Sync(); err != nil {
_ = tmp.Close()
return "", fmt.Errorf("sync temp file: %w", err)
}
if err := tmp.Close(); err != nil {
return "", fmt.Errorf("close temp file: %w", err)
}
if err := os.Chmod(tmpPath, perm); err != nil {
return "", fmt.Errorf("set temp file permissions: %w", err)
}
if err := os.Rename(tmpPath, dst); err != nil {
return "", fmt.Errorf("install temp file: %w", err)
}
removeTmp = false
return hex.EncodeToString(digest.Sum(nil)), nil
}
// InstallDownloadedTempFile installs a previously downloaded temp file at dst.
func InstallDownloadedTempFile(tmpPath, dst string, perm os.FileMode) error {
if strings.TrimSpace(tmpPath) == "" || strings.TrimSpace(dst) == "" {
return fmt.Errorf("temp and destination paths are required")
}
if err := os.MkdirAll(filepath.Dir(dst), 0o755); err != nil {
return fmt.Errorf("create destination directory: %w", err)
}
if err := os.Chmod(tmpPath, perm); err != nil {
return fmt.Errorf("set temp file permissions: %w", err)
}
if err := os.Rename(tmpPath, dst); err != nil {
return fmt.Errorf("install downloaded file: %w", err)
}
return nil
}

View File

@@ -0,0 +1,125 @@
package fileops
import (
"os"
"path/filepath"
"strings"
"testing"
)
func TestWriteFileAtomicOverwritesAndLeavesNoTempFile(t *testing.T) {
root := t.TempDir()
dst := filepath.Join(root, "out", "value.txt")
if err := WriteFileAtomic(dst, []byte("one"), 0o644); err != nil {
t.Fatalf("WriteFileAtomic(first) error = %v", err)
}
if err := WriteFileAtomic(dst, []byte("two"), 0o644); err != nil {
t.Fatalf("WriteFileAtomic(second) error = %v", err)
}
data, err := os.ReadFile(dst)
if err != nil {
t.Fatalf("ReadFile() error = %v", err)
}
if string(data) != "two" {
t.Fatalf("file content = %q, want %q", string(data), "two")
}
assertNoMatchingTempFiles(t, filepath.Dir(dst), "."+filepath.Base(dst)+".tmp-")
}
func TestWriteFileAtomicCleansTempFileOnInstallFailure(t *testing.T) {
root := t.TempDir()
blockedPath := filepath.Join(root, "blocked")
if err := os.MkdirAll(blockedPath, 0o755); err != nil {
t.Fatalf("MkdirAll(blockedPath) error = %v", err)
}
err := WriteFileAtomic(blockedPath, []byte("data"), 0o644)
if err == nil {
t.Fatal("WriteFileAtomic() error = nil, want install failure")
}
assertNoMatchingTempFiles(t, root, ".blocked.tmp-")
}
func TestCopyFileAtomicWithChecksumMatchesDestination(t *testing.T) {
root := t.TempDir()
src := filepath.Join(root, "source.txt")
dst := filepath.Join(root, "out", "copied.txt")
if err := os.WriteFile(src, []byte("copied-data"), 0o644); err != nil {
t.Fatalf("WriteFile(source) error = %v", err)
}
checksum, err := CopyFileAtomicWithChecksum(src, dst, 0o644)
if err != nil {
t.Fatalf("CopyFileAtomicWithChecksum() error = %v", err)
}
wantChecksum := "6e5c3f239e28cc315d57b2fcfc24169369c44a25802c0616a6d7081707fd24df"
if checksum != wantChecksum {
t.Fatalf("checksum = %q, want %q", checksum, wantChecksum)
}
data, err := os.ReadFile(dst)
if err != nil {
t.Fatalf("ReadFile(destination) error = %v", err)
}
if string(data) != "copied-data" {
t.Fatalf("destination content = %q, want %q", string(data), "copied-data")
}
assertNoMatchingTempFiles(t, filepath.Dir(dst), ".copied.txt.tmp-")
}
func TestCopyFileAtomicCleansTempFileOnInstallFailure(t *testing.T) {
root := t.TempDir()
src := filepath.Join(root, "source.txt")
if err := os.WriteFile(src, []byte("copied-data"), 0o644); err != nil {
t.Fatalf("WriteFile(source) error = %v", err)
}
blockedPath := filepath.Join(root, "blocked")
if err := os.MkdirAll(blockedPath, 0o755); err != nil {
t.Fatalf("MkdirAll(blockedPath) error = %v", err)
}
err := CopyFileAtomic(src, blockedPath, 0o644)
if err == nil {
t.Fatal("CopyFileAtomic() error = nil, want install failure")
}
assertNoMatchingTempFiles(t, root, ".blocked.tmp-")
}
func TestInstallDownloadedTempFileSetsPermissions(t *testing.T) {
root := t.TempDir()
tmpPath := filepath.Join(root, ".payload.tmp")
dst := filepath.Join(root, "out", "payload.json")
if err := os.WriteFile(tmpPath, []byte("{\"ok\":true}\n"), 0o600); err != nil {
t.Fatalf("WriteFile(temp) error = %v", err)
}
if err := InstallDownloadedTempFile(tmpPath, dst, 0o644); err != nil {
t.Fatalf("InstallDownloadedTempFile() error = %v", err)
}
if _, err := os.Stat(tmpPath); !os.IsNotExist(err) {
t.Fatalf("temp file still exists: stat err = %v", err)
}
info, err := os.Stat(dst)
if err != nil {
t.Fatalf("Stat(destination) error = %v", err)
}
if info.Mode().Perm() != 0o644 {
t.Fatalf("destination mode = %o, want 644", info.Mode().Perm())
}
}
func assertNoMatchingTempFiles(t *testing.T, dir, prefix string) {
t.Helper()
entries, err := os.ReadDir(dir)
if err != nil {
t.Fatalf("ReadDir(%q) error = %v", dir, err)
}
for _, e := range entries {
if strings.HasPrefix(e.Name(), prefix) {
t.Fatalf("unexpected temp file residue: %s", filepath.Join(dir, e.Name()))
}
}
}

View File

@@ -2,6 +2,8 @@ package pathsafe
import ( import (
"errors" "errors"
"path/filepath"
"strings"
"testing" "testing"
) )
@@ -41,3 +43,76 @@ func TestNormalizeRelativeDestination(t *testing.T) {
}) })
} }
} }
func TestJoinSlashRelativeUnderRoot(t *testing.T) {
root := filepath.Join(t.TempDir(), "session")
tests := []struct {
name string
input string
want string
wantErr error
}{
{name: "valid relative", input: "artifacts/session_recap.md", want: filepath.Join(root, "artifacts", "session_recap.md")},
{name: "windows separators normalized", input: `artifacts\session_recap.md`, want: filepath.Join(root, "artifacts", "session_recap.md")},
{name: "reject empty", input: "", wantErr: ErrRelativePathRequired},
{name: "reject absolute", input: "/artifacts/session_recap.md", wantErr: ErrRelativePathAbsolute},
{name: "reject traversal", input: "../artifacts/session_recap.md", wantErr: ErrRelativePathEscape},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
got, err := JoinSlashRelativeUnderRoot(root, tt.input)
if tt.wantErr != nil {
if !errors.Is(err, tt.wantErr) {
t.Fatalf("JoinSlashRelativeUnderRoot() error = %v, want %v", err, tt.wantErr)
}
return
}
if err != nil {
t.Fatalf("JoinSlashRelativeUnderRoot() error = %v", err)
}
if got != tt.want {
t.Fatalf("JoinSlashRelativeUnderRoot() = %q, want %q", got, tt.want)
}
})
}
}
func TestSlashRelativeFromRoot(t *testing.T) {
root := filepath.Join(t.TempDir(), "session")
target := filepath.Join(root, "transcripts", "full.json")
got, err := SlashRelativeFromRoot(root, target)
if err != nil {
t.Fatalf("SlashRelativeFromRoot() error = %v", err)
}
if got != "transcripts/full.json" {
t.Fatalf("SlashRelativeFromRoot() = %q, want transcripts/full.json", got)
}
got, err = SlashRelativeFromRoot(root, `transcripts\full.json`)
if err != nil {
t.Fatalf("SlashRelativeFromRoot(relative with windows separators) error = %v", err)
}
if got != "transcripts/full.json" {
t.Fatalf("SlashRelativeFromRoot(relative with windows separators) = %q, want transcripts/full.json", got)
}
}
func TestSlashRelativeFromRootRejectsOutsideRoot(t *testing.T) {
root := filepath.Join(t.TempDir(), "session")
outside := filepath.Join(filepath.Dir(root), "outside", "file.txt")
_, err := SlashRelativeFromRoot(root, outside)
if !errors.Is(err, ErrRelativePathEscape) {
t.Fatalf("SlashRelativeFromRoot() error = %v, want %v", err, ErrRelativePathEscape)
}
}
func TestJoinSlashRelativeUnderRootRequiresRoot(t *testing.T) {
_, err := JoinSlashRelativeUnderRoot("", "artifacts/session_recap.md")
if err == nil || !strings.Contains(err.Error(), "root path is required") {
t.Fatalf("JoinSlashRelativeUnderRoot() error = %v, want root-required error", err)
}
}

View File

@@ -0,0 +1,58 @@
package pathsafe
import (
"fmt"
"path/filepath"
"strings"
)
// JoinSlashRelativeUnderRoot validates a slash-style relative path and resolves
// it under root. The returned path uses the host filepath separator.
func JoinSlashRelativeUnderRoot(root, relative string) (string, error) {
rootClean := filepath.Clean(strings.TrimSpace(root))
if rootClean == "." || rootClean == "" {
return "", fmt.Errorf("root path is required")
}
normalized, err := NormalizeRelativeDestination(relative)
if err != nil {
return "", err
}
joined := filepath.Clean(filepath.Join(rootClean, filepath.FromSlash(normalized)))
rel, err := filepath.Rel(rootClean, joined)
if err != nil {
return "", fmt.Errorf("resolve relative path under root: %w", err)
}
if rel == "." || rel == ".." || strings.HasPrefix(rel, ".."+string(filepath.Separator)) {
return "", ErrRelativePathEscape
}
return joined, nil
}
// SlashRelativeFromRoot derives a slash-style relative path for target under
// root. Target may be absolute or relative to root.
func SlashRelativeFromRoot(root, target string) (string, error) {
rootClean := filepath.Clean(strings.TrimSpace(root))
if rootClean == "." || rootClean == "" {
return "", fmt.Errorf("root path is required")
}
targetClean := filepath.Clean(strings.TrimSpace(target))
if targetClean == "." || targetClean == "" {
return "", ErrRelativePathRequired
}
if !filepath.IsAbs(targetClean) {
targetClean = filepath.Clean(filepath.Join(rootClean, targetClean))
}
rel, err := filepath.Rel(rootClean, targetClean)
if err != nil {
return "", fmt.Errorf("derive path relative to root: %w", err)
}
normalized, err := NormalizeRelativeDestination(filepath.ToSlash(rel))
if err != nil {
return "", err
}
return normalized, nil
}

View File

@@ -229,7 +229,11 @@ func artifactRelativePathCandidates(
candidates = append(candidates, normalized) candidates = append(candidates, normalized)
} }
sourceDescriptor, err := artifactpolicy.PreviousSessionSourceDescriptorForConfiguredKey(artifactName)
sourceID := artifactpolicy.ConfiguredSourceID(artifactName) sourceID := artifactpolicy.ConfiguredSourceID(artifactName)
if err == nil {
sourceID = sourceDescriptor.ConfiguredSourceID
}
if rel, ok := manifestArtifactRelativePathBySourceID(previousManifest, sourceID); ok { if rel, ok := manifestArtifactRelativePathBySourceID(previousManifest, sourceID); ok {
appendCandidate(rel) appendCandidate(rel)
base := path.Base(rel) base := path.Base(rel)
@@ -309,11 +313,7 @@ func deriveManifestRelativePath(previousManifest *manifest.Manifest, localPath s
if !ok { if !ok {
return "", false return "", false
} }
rel, err := filepath.Rel(sessionRoot, trimmed) normalized, err := pathsafe.SlashRelativeFromRoot(sessionRoot, trimmed)
if err != nil {
return "", false
}
normalized, err := pathsafe.NormalizeRelativeDestination(filepath.ToSlash(rel))
if err != nil { if err != nil {
return "", false return "", false
} }
@@ -375,11 +375,7 @@ func relativeToSession(paths artifacts.SessionPaths, localPath string) (string,
if strings.TrimSpace(root) == "" { if strings.TrimSpace(root) == "" {
return "", fmt.Errorf("session root is required") return "", fmt.Errorf("session root is required")
} }
rel, err := filepath.Rel(root, filepath.Clean(localPath)) normalized, err := pathsafe.SlashRelativeFromRoot(root, localPath)
if err != nil {
return "", fmt.Errorf("resolve previous-cache relative path: %w", err)
}
normalized, err := pathsafe.NormalizeRelativeDestination(filepath.ToSlash(rel))
if err != nil { if err != nil {
return "", fmt.Errorf("resolve previous-cache relative path: %w", err) return "", fmt.Errorf("resolve previous-cache relative path: %w", err)
} }

View File

@@ -628,11 +628,11 @@ func resolveScriptoriumInput(
runtimeCatalog *artifacts.ArtifactCatalog, runtimeCatalog *artifacts.ArtifactCatalog,
) (string, bool, *artifacts.ResolvedSessionArtifact, error) { ) (string, bool, *artifacts.ResolvedSessionArtifact, error) {
source := strings.TrimSpace(inputCfg.Source) source := strings.TrimSpace(inputCfg.Source)
classified, classifyErr := artifactpolicy.ClassifySource(source) descriptor, describeErr := artifactpolicy.DescribeScriptoriumInputSource(source)
if classifyErr != nil { if describeErr != nil {
return "", false, nil, classifyErr return "", false, nil, describeErr
} }
if classified.Kind == artifactpolicy.SourceKindPreviousArtifact { if descriptor.Source.Kind == artifactpolicy.SourceKindPreviousArtifact {
resolved, err := artifacts.ResolvePreviousSessionArtifactWithCatalog(paths, m, source, runtimeCatalog) resolved, err := artifacts.ResolvePreviousSessionArtifactWithCatalog(paths, m, source, runtimeCatalog)
if err == nil { if err == nil {
copy := resolved copy := resolved
@@ -657,13 +657,13 @@ func resolveScriptoriumInput(
return resolved.Path, true, &copy, nil return resolved.Path, true, &copy, nil
} }
if errors.Is(err, artifacts.ErrSessionArtifactNotFound) { if errors.Is(err, artifacts.ErrSessionArtifactNotFound) {
if classified.Kind == artifactpolicy.SourceKindConfiguredArtifact { if descriptor.Source.Kind == artifactpolicy.SourceKindConfiguredArtifact {
if inputCfg.Required { if inputCfg.Required {
return "", false, nil, fmt.Errorf("configured artifact source %q is unavailable", source) return "", false, nil, fmt.Errorf("configured artifact source %q is unavailable", source)
} }
return "", false, nil, nil return "", false, nil, nil
} }
switch classified.ID { switch descriptor.Source.ID {
case artifacts.ArtifactTranscriptPolished: case artifacts.ArtifactTranscriptPolished:
return "", false, nil, nil return "", false, nil, nil
case artifacts.ArtifactTranscriptFinal: case artifacts.ArtifactTranscriptFinal:

View File

@@ -8,6 +8,7 @@ import (
"sort" "sort"
"gitea.maximumdirect.net/eric/narratio/internal/artifacts" "gitea.maximumdirect.net/eric/narratio/internal/artifacts"
"gitea.maximumdirect.net/eric/narratio/internal/fileops"
"gitea.maximumdirect.net/eric/narratio/internal/manifest" "gitea.maximumdirect.net/eric/narratio/internal/manifest"
"gitea.maximumdirect.net/eric/narratio/internal/previouscache" "gitea.maximumdirect.net/eric/narratio/internal/previouscache"
) )
@@ -54,8 +55,35 @@ func hydratePreviousSessionArtifacts(
if err := os.MkdirAll(filepath.Dir(record.LocalPath), 0o755); err != nil { if err := os.MkdirAll(filepath.Dir(record.LocalPath), 0o755); err != nil {
return nil, fmt.Errorf("create previous-session path directory for %q: %w", record.LocalPath, err) return nil, fmt.Errorf("create previous-session path directory for %q: %w", record.LocalPath, err)
} }
if err := env.ObjectStore.Download(ctx, record.RemoteKey, record.LocalPath); err != nil {
return nil, fmt.Errorf("download previous-session object %q to %q: %w", record.RemoteKey, record.LocalPath, err) base := filepath.Base(record.LocalPath)
tmp, err := os.CreateTemp(filepath.Dir(record.LocalPath), "."+base+".prepare-previous-*.tmp")
if err != nil {
return nil, fmt.Errorf("create previous-session temp file for %q: %w", record.LocalPath, err)
}
tmpPath := tmp.Name()
if err := tmp.Close(); err != nil {
_ = os.Remove(tmpPath)
return nil, fmt.Errorf("close previous-session temp file for %q: %w", record.LocalPath, err)
}
if err := func() error {
removeTmp := true
defer func() {
if removeTmp {
_ = os.Remove(tmpPath)
}
}()
if err := env.ObjectStore.Download(ctx, record.RemoteKey, tmpPath); err != nil {
return fmt.Errorf("download previous-session object %q to temp file: %w", record.RemoteKey, err)
}
if err := fileops.InstallDownloadedTempFile(tmpPath, record.LocalPath, 0o644); err != nil {
return fmt.Errorf("install previous-session object %q at %q: %w", record.RemoteKey, record.LocalPath, err)
}
removeTmp = false
return nil
}(); err != nil {
return nil, err
} }
if record.Kind == preparePreviousInputKindArtifact { if record.Kind == preparePreviousInputKindArtifact {
if err := requireNonEmptyFile(record.LocalPath, "previous-session artifact "+record.RequirementName); err != nil { if err := requireNonEmptyFile(record.LocalPath, "previous-session artifact "+record.RequirementName); err != nil {

View File

@@ -9,6 +9,7 @@ import (
"gitea.maximumdirect.net/eric/narratio/internal/artifacts" "gitea.maximumdirect.net/eric/narratio/internal/artifacts"
"gitea.maximumdirect.net/eric/narratio/internal/config" "gitea.maximumdirect.net/eric/narratio/internal/config"
"gitea.maximumdirect.net/eric/narratio/internal/manifest" "gitea.maximumdirect.net/eric/narratio/internal/manifest"
"gitea.maximumdirect.net/eric/narratio/internal/pathsafe"
) )
type runStageLayout struct { type runStageLayout struct {
@@ -92,18 +93,17 @@ func runLocalPathForCanonical(layout runStageLayout, sessionPaths artifacts.Sess
if cleanCanonical == "" { if cleanCanonical == "" {
return "", fmt.Errorf("canonical path is required") return "", fmt.Errorf("canonical path is required")
} }
rel, err := filepath.Rel(filepath.Clean(sessionPaths.Root), cleanCanonical) rel, err := pathsafe.SlashRelativeFromRoot(sessionPaths.Root, cleanCanonical)
if err != nil { if err != nil {
return "", fmt.Errorf("derive session-relative path for %q: %w", cleanCanonical, err) return "", fmt.Errorf("derive session-relative path for %q: %w", cleanCanonical, err)
} }
rel = filepath.Clean(rel) if rel == config.PathPreviousDirSegment || strings.HasPrefix(rel, config.PathPreviousDirSegment+"/") {
if rel == "." || rel == ".." || strings.HasPrefix(rel, ".."+string(filepath.Separator)) {
return "", fmt.Errorf("canonical path %q is outside session root %q", cleanCanonical, sessionPaths.Root)
}
if rel == config.PathPreviousDirSegment || strings.HasPrefix(rel, config.PathPreviousDirSegment+string(filepath.Separator)) {
return cleanCanonical, nil return cleanCanonical, nil
} }
localPath := filepath.Join(layout.OutputsDir, rel) localPath, err := pathsafe.JoinSlashRelativeUnderRoot(layout.OutputsDir, rel)
if err != nil {
return "", fmt.Errorf("resolve run-local output path for %q: %w", cleanCanonical, err)
}
if err := os.MkdirAll(filepath.Dir(localPath), 0o755); err != nil { if err := os.MkdirAll(filepath.Dir(localPath), 0o755); err != nil {
return "", fmt.Errorf("create run-local output parent for %q: %w", localPath, err) return "", fmt.Errorf("create run-local output parent for %q: %w", localPath, err)
} }