From 2ca700195c1c9ede820a1c599d29052be5723e43 Mon Sep 17 00:00:00 2001 From: Eric Rakestraw Date: Wed, 20 May 2026 14:53:22 +0000 Subject: [PATCH] Integrate previous-session artifact hydration into prepare stage --- internal/app/runner.go | 13 +++ internal/app/runner_test.go | 49 +++++++++ internal/stage/prepare.go | 69 +++++++++++-- internal/stage/prepare_test.go | 182 +++++++++++++++++++++++++++++++++ 4 files changed, 307 insertions(+), 6 deletions(-) diff --git a/internal/app/runner.go b/internal/app/runner.go index 77c35d3..c8f3f4d 100644 --- a/internal/app/runner.go +++ b/internal/app/runner.go @@ -555,6 +555,12 @@ func needsObjectStoreForRun(cfg *config.Config, stages []stage.Stage) bool { if cfg.Session.Inputs.AudioS3 != nil && stageRequested("prepare") { return true } + if stageRequested("prepare") { + requirements := artifacts.CollectPreviousArtifactRequirements(configuredScriptoriumArtifacts(cfg)) + if len(requirements) > 0 && strings.TrimSpace(cfg.Session.PreviousSessionID) != "" { + return true + } + } if !stageRequested("archive") { return false } @@ -569,3 +575,10 @@ func needsObjectStoreForRun(cfg *config.Config, stages []stage.Stage) bool { } return true } + +func configuredScriptoriumArtifacts(cfg *config.Config) map[string]config.ScriptoriumArtifactConfig { + if cfg == nil || cfg.Pipeline == nil || cfg.Pipeline.Scriptorium == nil { + return nil + } + return cfg.Pipeline.Scriptorium.Artifacts +} diff --git a/internal/app/runner_test.go b/internal/app/runner_test.go index 6278909..0278138 100644 --- a/internal/app/runner_test.go +++ b/internal/app/runner_test.go @@ -199,6 +199,55 @@ func TestExecuteStagesAnalyzeOutputsPersistAsScriptoriumArtifacts(t *testing.T) } } +func TestNeedsObjectStoreForRunPrepareWithPreviousRequirements(t *testing.T) { + tests := []struct { + name string + previousSessionID string + want bool + }{ + { + name: "previous session configured", + previousSessionID: "2026-05-10", + want: true, + }, + { + name: "previous session missing", + previousSessionID: "", + want: false, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + cfg := &config.Config{ + Pipeline: &config.PipelineConfig{ + Scriptorium: &config.ScriptoriumConfig{ + Artifacts: map[string]config.ScriptoriumArtifactConfig{ + "session_recap": { + Enabled: true, + Inputs: map[string]config.ScriptoriumInputConfig{ + "previous_recap": { + Source: "narratio.previous_session.artifact.session_recap", + Required: true, + }, + }, + }, + }, + }, + }, + Session: &config.SessionConfig{ + PreviousSessionID: tt.previousSessionID, + }, + } + + got := needsObjectStoreForRun(cfg, []stage.Stage{countingStage{name: "prepare", runs: new(int)}}) + if got != tt.want { + t.Fatalf("needsObjectStoreForRun() = %v, want %v", got, tt.want) + } + }) + } +} + func TestExecuteStagesArchiveFailsWhenRequiredRecapPromotionMissingForSelectedArtifacts(t *testing.T) { cfg := testConfig(t) cfg.Pipeline.Storage.S3 = &config.StorageS3Config{ diff --git a/internal/stage/prepare.go b/internal/stage/prepare.go index 6526196..45a4a8d 100644 --- a/internal/stage/prepare.go +++ b/internal/stage/prepare.go @@ -145,6 +145,20 @@ func (prepareStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*S } } + previousRequirements := collectPreparePreviousRequirements(env.Config) + var previousHydration *previousSessionHydrationResult + if len(previousRequirements) > 0 { + if err := clearManagedPreviousState(paths); err != nil { + return nil, fmt.Errorf("prepare: clear previous-session cache: %w", err) + } + hydration, err := hydratePreviousSessionArtifacts(ctx, env, paths, previousRequirements) + if err != nil { + return nil, fmt.Errorf("prepare: hydrate previous-session artifacts: %w", err) + } + previousHydration = hydration + inputs = append(inputs, hydration.Inputs...) + } + sort.Slice(inputs, func(i, j int) bool { if inputs[i].Kind != inputs[j].Kind { return inputs[i].Kind < inputs[j].Kind @@ -153,13 +167,27 @@ func (prepareStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*S }) m.Inputs = inputs + metadata := map[string]any{ + "prepared": true, + "stage": "prepare", + "inputs_count": len(inputs), + "audio_files_resolved": countAudioInputs(inputs), + } + if len(previousRequirements) > 0 { + metadata["previous_requirements_count"] = len(previousRequirements) + if previousHydration != nil { + metadata["previous_artifacts_hydrated"] = append([]string(nil), previousHydration.Hydrated...) + metadata["previous_artifacts_hydrated_count"] = len(previousHydration.Hydrated) + metadata["previous_artifacts_missing_optional"] = append([]string(nil), previousHydration.SkippedMissing...) + metadata["previous_artifacts_missing_optional_count"] = len(previousHydration.SkippedMissing) + if strings.TrimSpace(previousHydration.PreviousRunID) != "" { + metadata["previous_session_run_id"] = strings.TrimSpace(previousHydration.PreviousRunID) + } + } + } + return &StageResult{ - Metadata: map[string]any{ - "prepared": true, - "stage": "prepare", - "inputs_count": len(inputs), - "audio_files_resolved": countAudioInputs(inputs), - }, + Metadata: metadata, }, nil } @@ -354,6 +382,35 @@ func countAudioInputs(inputs []manifest.InputRecord) int { return count } +func collectPreparePreviousRequirements(cfg *config.Config) []artifacts.PreviousArtifactRequirement { + if cfg == nil || cfg.Pipeline == nil || cfg.Pipeline.Scriptorium == nil { + return nil + } + return artifacts.CollectPreviousArtifactRequirements(cfg.Pipeline.Scriptorium.Artifacts) +} + +func clearManagedPreviousState(paths artifacts.SessionPaths) error { + previousDir := filepath.Clean(paths.PreviousDir) + sessionRoot := filepath.Clean(paths.Root) + if strings.TrimSpace(previousDir) == "" || strings.TrimSpace(sessionRoot) == "" { + return fmt.Errorf("previous/session root paths are required") + } + if previousDir == sessionRoot { + return fmt.Errorf("refusing to clear session root as previous cache: %q", previousDir) + } + prefix := sessionRoot + string(filepath.Separator) + if !strings.HasPrefix(previousDir, prefix) { + return fmt.Errorf("refusing to clear path outside session root: %q", previousDir) + } + if filepath.Base(previousDir) != config.PathPreviousDirSegment { + return fmt.Errorf("refusing to clear non-previous path %q", previousDir) + } + if err := os.RemoveAll(previousDir); err != nil { + return err + } + return os.MkdirAll(previousDir, 0o755) +} + func pathsWorkDirForManifest(env *Env, m *manifest.Manifest, sessionID string) string { if env == nil || env.Config == nil || env.Config.Pipeline == nil || env.Config.Session == nil { return "" diff --git a/internal/stage/prepare_test.go b/internal/stage/prepare_test.go index 0209e4c..f5cef4b 100644 --- a/internal/stage/prepare_test.go +++ b/internal/stage/prepare_test.go @@ -284,6 +284,188 @@ func TestPrepareStageAudioSourceConflictFails(t *testing.T) { } } +func TestPrepareStageWithoutPreviousRequirementsDoesNotTouchPreviousState(t *testing.T) { + env, m := setupPrepareEnv(t) + root := filepath.Dir(env.Config.SessionPath) + writeFile(t, filepath.Join(root, "audio", "a.flac"), "a") + env.Config.Session.Inputs.AudioFiles = []string{"./audio/a.flac"} + + paths := sessionPathsForEnv(env, m.SessionID) + stalePath := filepath.Join(paths.PreviousDir, "stale.txt") + writeFile(t, stalePath, "stale") + + _, err := (prepareStage{}).Run(context.Background(), env, m) + if err != nil { + t.Fatalf("prepare.Run() error = %v", err) + } + if _, err := os.Stat(stalePath); err != nil { + t.Fatalf("expected previous stale file to remain untouched: %v", err) + } +} + +func TestPrepareStageOptionalPreviousArtifactWithoutPreviousSessionIDSucceeds(t *testing.T) { + env, m := setupPrepareEnv(t) + root := filepath.Dir(env.Config.SessionPath) + writeFile(t, filepath.Join(root, "audio", "a.flac"), "a") + env.Config.Session.Inputs.AudioFiles = []string{"./audio/a.flac"} + configurePreparePreviousArtifactSource(env, false) + env.Config.Session.PreviousSessionID = "" + + paths := sessionPathsForEnv(env, m.SessionID) + writeFile(t, filepath.Join(paths.PreviousArtifactsDir, "stale.txt"), "stale") + + result, err := (prepareStage{}).Run(context.Background(), env, m) + if err != nil { + t.Fatalf("prepare.Run() error = %v", err) + } + if result.Metadata["previous_requirements_count"] != 1 { + t.Fatalf("metadata previous_requirements_count = %#v, want 1", result.Metadata["previous_requirements_count"]) + } + if result.Metadata["previous_artifacts_hydrated_count"] != 0 { + t.Fatalf("metadata previous_artifacts_hydrated_count = %#v, want 0", result.Metadata["previous_artifacts_hydrated_count"]) + } + if result.Metadata["previous_artifacts_missing_optional_count"] != 1 { + t.Fatalf("metadata previous_artifacts_missing_optional_count = %#v, want 1", result.Metadata["previous_artifacts_missing_optional_count"]) + } + if _, err := os.Stat(filepath.Join(paths.PreviousArtifactsDir, "stale.txt")); !os.IsNotExist(err) { + t.Fatalf("expected stale previous state to be cleared, stat err = %v", err) + } + for _, in := range m.Inputs { + if in.Kind == preparePreviousInputKindManifest || in.Kind == preparePreviousInputKindArtifact { + t.Fatalf("unexpected previous input record: %#v", in) + } + } +} + +func TestPrepareStageRequiredPreviousArtifactWithoutPreviousSessionIDFails(t *testing.T) { + env, m := setupPrepareEnv(t) + root := filepath.Dir(env.Config.SessionPath) + writeFile(t, filepath.Join(root, "audio", "a.flac"), "a") + env.Config.Session.Inputs.AudioFiles = []string{"./audio/a.flac"} + configurePreparePreviousArtifactSource(env, true) + env.Config.Session.PreviousSessionID = "" + + _, err := (prepareStage{}).Run(context.Background(), env, m) + if err == nil || !strings.Contains(err.Error(), "previous_session_id is required") { + t.Fatalf("error = %v, want previous_session_id required failure", err) + } +} + +func TestPrepareStageHydratesRequiredPreviousArtifactAndRecordsInputs(t *testing.T) { + env, m := setupPrepareEnv(t) + root := filepath.Dir(env.Config.SessionPath) + writeFile(t, filepath.Join(root, "audio", "a.flac"), "a") + env.Config.Session.Inputs.AudioFiles = []string{"./audio/a.flac"} + env.Config.Session.Campaign = "forsaken" + env.Config.Session.PreviousSessionID = "2026-05-10" + env.Config.Pipeline.Storage.S3 = &config.StorageS3Config{Bucket: "my-dnd-archive", RootPrefix: "dnd"} + configurePreparePreviousArtifactSource(env, true) + + fake := &storage.FakeBackend{} + env.ObjectStore = fake + seedPreviousCurrentState(t, env, fake, previousStateSeedOptions{ + includeRunPointerObject: true, + includeArtifactObject: true, + artifactBody: "# prior recap\n", + }) + + paths := sessionPathsForEnv(env, m.SessionID) + result, err := (prepareStage{}).Run(context.Background(), env, m) + if err != nil { + t.Fatalf("prepare.Run() error = %v", err) + } + + if _, err := os.Stat(paths.PreviousManifestPath); err != nil { + t.Fatalf("expected previous manifest: %v", err) + } + recapPath := artifacts.SessionPreviousArtifactPath(paths, "artifacts/session_recap.md") + if _, err := os.Stat(recapPath); err != nil { + t.Fatalf("expected previous artifact: %v", err) + } + + var hasPreviousManifest, hasPreviousArtifact bool + for _, in := range m.Inputs { + if in.Kind == preparePreviousInputKindManifest { + hasPreviousManifest = true + } + if in.Kind == preparePreviousInputKindArtifact { + hasPreviousArtifact = true + } + } + if !hasPreviousManifest || !hasPreviousArtifact { + t.Fatalf("manifest inputs missing previous provenance: %#v", m.Inputs) + } + if result.Metadata["previous_artifacts_hydrated_count"] != 1 { + t.Fatalf("metadata previous_artifacts_hydrated_count = %#v, want 1", result.Metadata["previous_artifacts_hydrated_count"]) + } +} + +func TestPrepareStageRerunOverwritesPreviousCache(t *testing.T) { + env, m := setupPrepareEnv(t) + root := filepath.Dir(env.Config.SessionPath) + writeFile(t, filepath.Join(root, "audio", "a.flac"), "a") + env.Config.Session.Inputs.AudioFiles = []string{"./audio/a.flac"} + env.Config.Session.Campaign = "forsaken" + env.Config.Session.PreviousSessionID = "2026-05-10" + env.Config.Pipeline.Storage.S3 = &config.StorageS3Config{Bucket: "my-dnd-archive", RootPrefix: "dnd"} + configurePreparePreviousArtifactSource(env, true) + + fake := &storage.FakeBackend{} + env.ObjectStore = fake + seed := seedPreviousCurrentState(t, env, fake, previousStateSeedOptions{ + includeRunPointerObject: true, + includeArtifactObject: true, + artifactBody: "# old recap\n", + }) + + paths := sessionPathsForEnv(env, m.SessionID) + if _, err := (prepareStage{}).Run(context.Background(), env, m); err != nil { + t.Fatalf("first prepare run error = %v", err) + } + recapPath := artifacts.SessionPreviousArtifactPath(paths, "artifacts/session_recap.md") + firstBytes, err := os.ReadFile(recapPath) + if err != nil { + t.Fatalf("read first hydrated artifact: %v", err) + } + if strings.TrimSpace(string(firstBytes)) != "# old recap" { + t.Fatalf("first hydrated content = %q, want %q", strings.TrimSpace(string(firstBytes)), "# old recap") + } + + writeFile(t, filepath.Join(paths.PreviousArtifactsDir, "stale.txt"), "stale") + fake.SeedObject(storage.FakeObject{Key: seed.ArtifactKey, Data: []byte("# new recap\n")}) + + if _, err := (prepareStage{}).Run(context.Background(), env, m); err != nil { + t.Fatalf("second prepare run error = %v", err) + } + secondBytes, err := os.ReadFile(recapPath) + if err != nil { + t.Fatalf("read second hydrated artifact: %v", err) + } + if strings.TrimSpace(string(secondBytes)) != "# new recap" { + t.Fatalf("second hydrated content = %q, want %q", strings.TrimSpace(string(secondBytes)), "# new recap") + } + if _, err := os.Stat(filepath.Join(paths.PreviousArtifactsDir, "stale.txt")); !os.IsNotExist(err) { + t.Fatalf("expected stale previous cache file to be removed, stat err = %v", err) + } +} + +func configurePreparePreviousArtifactSource(env *Env, required bool) { + env.Config.Pipeline.Scriptorium = &config.ScriptoriumConfig{ + Artifacts: map[string]config.ScriptoriumArtifactConfig{ + "session_recap": { + Enabled: true, + OutputPath: "artifacts/session_recap.md", + Inputs: map[string]config.ScriptoriumInputConfig{ + "previous_recap": { + Source: "narratio.previous_session.artifact.session_recap", + Required: required, + }, + }, + }, + }, + } +} + func setupPrepareEnv(t *testing.T) (*Env, *manifest.Manifest) { t.Helper() workspace := t.TempDir()