package app import ( "context" "errors" "os" "path/filepath" "strings" "testing" "time" "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/manifest" "gitea.maximumdirect.net/eric/narratio/internal/stage" ) type publishSuccessStage struct { metadata map[string]any } func (publishSuccessStage) Name() string { return "publish" } func (publishSuccessStage) Declares() stage.IODecl { return stage.IODecl{} } func (s publishSuccessStage) Run(_ context.Context, _ *stage.Env, _ *manifest.Manifest) (*stage.StageResult, error) { md := map[string]any{ "stage": "publish", "uploaded": true, "current_pointer_written": true, "current_run_id_key": "dnd/campaigns/sample-campaign/sessions/2026-05-03/current/run_id.txt", } for k, v := range s.metadata { md[k] = v } return &stage.StageResult{Metadata: md}, nil } type notifyFailStage struct{} func (notifyFailStage) Name() string { return "notify" } func (notifyFailStage) Declares() stage.IODecl { return stage.IODecl{} } func (notifyFailStage) Run(_ context.Context, _ *stage.Env, _ *manifest.Manifest) (*stage.StageResult, error) { return nil, errors.New("notify failed") } func TestPostPublishCleanupDisabledKeepsLocalDirs(t *testing.T) { cfg, seed := cleanupFixtureConfig(t) cfg.Pipeline.Spool.DeleteAudioAfterPublish = false cfg.Pipeline.Workspace.CleanupAfterPublish = false if _, err := executeStages(context.Background(), cfg, []stage.Stage{publishSuccessStage{}}, RunOptions{Env: &Env{ObjectStore: &storage.FakeBackend{}}}); err != nil { t.Fatalf("executeStages() error = %v", err) } assertExists(t, seed.spoolAudioDir) assertExists(t, seed.runWorkDir) assertExists(t, seed.localSourceAudio) } func TestPostPublishCleanupSpoolOnly(t *testing.T) { cfg, seed := cleanupFixtureConfig(t) cfg.Pipeline.Spool.DeleteAudioAfterPublish = true cfg.Pipeline.Workspace.CleanupAfterPublish = false if _, err := executeStages(context.Background(), cfg, []stage.Stage{publishSuccessStage{}}, RunOptions{Env: &Env{ObjectStore: &storage.FakeBackend{}}}); err != nil { t.Fatalf("executeStages() error = %v", err) } assertMissing(t, seed.spoolAudioDir) assertExists(t, seed.runWorkDir) assertExists(t, seed.localSourceAudio) } func TestPostPublishCleanupWorkdirOnly(t *testing.T) { cfg, seed := cleanupFixtureConfig(t) cfg.Pipeline.Spool.DeleteAudioAfterPublish = false cfg.Pipeline.Workspace.CleanupAfterPublish = true if _, err := executeStages(context.Background(), cfg, []stage.Stage{publishSuccessStage{}}, RunOptions{Env: &Env{ObjectStore: &storage.FakeBackend{}}}); err != nil { t.Fatalf("executeStages() error = %v", err) } assertExists(t, cfg.Pipeline.Workspace.Root) assertExists(t, seed.otherRunDir) assertExists(t, seed.previousCachePath) assertMissing(t, seed.runWorkDir) assertExists(t, seed.spoolAudioDir) } func TestPostPublishCleanupBothPolicies(t *testing.T) { cfg, seed := cleanupFixtureConfig(t) cfg.Pipeline.Spool.DeleteAudioAfterPublish = true cfg.Pipeline.Workspace.CleanupAfterPublish = true if _, err := executeStages(context.Background(), cfg, []stage.Stage{publishSuccessStage{}}, RunOptions{Env: &Env{ObjectStore: &storage.FakeBackend{}}}); err != nil { t.Fatalf("executeStages() error = %v", err) } assertMissing(t, seed.spoolAudioDir) assertMissing(t, seed.runWorkDir) assertExists(t, seed.otherRunDir) assertExists(t, seed.previousCachePath) } func TestPostPublishCleanupNotRunWhenPublishFails(t *testing.T) { cfg, seed := cleanupFixtureConfig(t) cfg.Pipeline.Spool.DeleteAudioAfterPublish = true cfg.Pipeline.Workspace.CleanupAfterPublish = true _, err := executeStages(context.Background(), cfg, []stage.Stage{failingStage{name: "publish", err: errors.New("publish failed")}}, RunOptions{Env: &Env{ObjectStore: &storage.FakeBackend{}}}) if err == nil || !strings.Contains(err.Error(), "stage \"publish\" failed") { t.Fatalf("executeStages() error = %v, want publish failure", err) } assertExists(t, seed.spoolAudioDir) assertExists(t, seed.runWorkDir) } func TestPostPublishCleanupNotRunWhenPublishSkipped(t *testing.T) { cfg, seed := cleanupFixtureConfig(t) cfg.Pipeline.Spool.DeleteAudioAfterPublish = true cfg.Pipeline.Workspace.CleanupAfterPublish = true if _, err := executeStages(context.Background(), cfg, []stage.Stage{publishSuccessStage{metadata: map[string]any{"skipped": true}}}, RunOptions{Env: &Env{ObjectStore: &storage.FakeBackend{}}}); err != nil { t.Fatalf("executeStages() error = %v", err) } assertExists(t, seed.spoolAudioDir) assertExists(t, seed.runWorkDir) } func TestPostPublishCleanupNotRunWhenCurrentPointerMissing(t *testing.T) { cfg, seed := cleanupFixtureConfig(t) cfg.Pipeline.Spool.DeleteAudioAfterPublish = true cfg.Pipeline.Workspace.CleanupAfterPublish = true if _, err := executeStages(context.Background(), cfg, []stage.Stage{publishSuccessStage{metadata: map[string]any{"current_pointer_written": false}}}, RunOptions{Env: &Env{ObjectStore: &storage.FakeBackend{}}}); err != nil { t.Fatalf("executeStages() error = %v", err) } assertExists(t, seed.spoolAudioDir) assertExists(t, seed.runWorkDir) } func TestPostPublishCleanupNotRunWhenPublishUploadDisabled(t *testing.T) { cfg, seed := cleanupFixtureConfig(t) cfg.Pipeline.Spool.DeleteAudioAfterPublish = true cfg.Pipeline.Workspace.CleanupAfterPublish = true cfg.Pipeline.Publish.UploadRun = boolPtr(false) if _, err := executeStages(context.Background(), cfg, []stage.Stage{publishSuccessStage{}}, RunOptions{Env: &Env{ObjectStore: &storage.FakeBackend{}}}); err != nil { t.Fatalf("executeStages() error = %v", err) } assertExists(t, seed.spoolAudioDir) assertExists(t, seed.runWorkDir) } func TestPostPublishCleanupWaitsUntilAllStagesSucceed(t *testing.T) { cfg, seed := cleanupFixtureConfig(t) cfg.Pipeline.Spool.DeleteAudioAfterPublish = true cfg.Pipeline.Workspace.CleanupAfterPublish = true _, err := executeStages(context.Background(), cfg, []stage.Stage{publishSuccessStage{}, notifyFailStage{}}, RunOptions{Env: &Env{ObjectStore: &storage.FakeBackend{}}}) if err == nil || !strings.Contains(err.Error(), "stage \"notify\" failed") { t.Fatalf("executeStages() error = %v, want notify failure", err) } assertExists(t, seed.spoolAudioDir) assertExists(t, seed.runWorkDir) } func TestPostPublishCleanupFailsOnUnsafePath(t *testing.T) { cfg, _ := cleanupFixtureConfig(t) cfg.Pipeline.Spool.DeleteAudioAfterPublish = true cfg.Pipeline.Workspace.CleanupAfterPublish = false manifestPath := manifestPathFor(cfg) store := &manifest.LocalStore{} m, err := store.Load(context.Background(), manifestPath) if err != nil { t.Fatalf("Load() error = %v", err) } m.LocalSpoolDir = filepath.Join(filepath.Dir(cfg.Pipeline.Spool.Root), "outside-spool") if err := store.Save(context.Background(), manifestPath, m); err != nil { t.Fatalf("Save() error = %v", err) } _, err = executeStages(context.Background(), cfg, []stage.Stage{publishSuccessStage{}}, RunOptions{Env: &Env{ObjectStore: &storage.FakeBackend{}}}) if err == nil || !strings.Contains(err.Error(), "refusing to delete path outside root") { t.Fatalf("executeStages() error = %v, want safe-path failure", err) } } func TestPostPublishCleanupNotRunWhenOutputIsMissing(t *testing.T) { cfg, seed, runID := publishStageCleanupFixture(t) cfg.Pipeline.Spool.DeleteAudioAfterPublish = true cfg.Pipeline.Workspace.CleanupAfterPublish = true cfg.Pipeline.Publish.Outputs = []config.PublishOutputRule{ {Source: "narratio.transcript.base", Dest: "transcripts/base.json", Required: boolPtr(true)}, } publishStageImpl, err := stage.Select("publish") if err != nil { t.Fatalf("Select(publish) error = %v", err) } _, err = executeStages(context.Background(), cfg, []stage.Stage{publishStageImpl}, RunOptions{Env: &Env{ObjectStore: &storage.FakeBackend{}}}) if err == nil || !strings.Contains(err.Error(), "required output source unavailable") { t.Fatalf("executeStages() error = %v, want required output source unavailable failure", err) } assertExists(t, seed.spoolAudioDir) assertExists(t, seed.runWorkDir) assertExists(t, filepath.Join(seed.runWorkDir, "manifest.json")) assertExists(t, artifacts.SessionRunRootForCampaign(cfg.Pipeline.Workspace.Root, cfg.Session.Campaign, cfg.Session.SessionID, runID)) } func TestPostPublishCleanupNotRunWhenCurrentManifestUploadFails(t *testing.T) { cfg, seed, _ := publishStageCleanupFixture(t) cfg.Pipeline.Spool.DeleteAudioAfterPublish = true cfg.Pipeline.Workspace.CleanupAfterPublish = true failKey := seed.sessionPrefix + "current/manifest.json" publishStageImpl, err := stage.Select("publish") if err != nil { t.Fatalf("Select(publish) error = %v", err) } _, err = executeStages(context.Background(), cfg, []stage.Stage{publishStageImpl}, RunOptions{ Env: &Env{ObjectStore: &failKeyStore{delegate: &storage.FakeBackend{}, failKey: failKey}}, }) if err == nil || !strings.Contains(err.Error(), "current manifest") { t.Fatalf("executeStages() error = %v, want current-manifest failure", err) } assertExists(t, seed.spoolAudioDir) assertExists(t, seed.runWorkDir) } func TestPostPublishCleanupNotRunWhenCurrentPointerUploadFails(t *testing.T) { cfg, seed, _ := publishStageCleanupFixture(t) cfg.Pipeline.Spool.DeleteAudioAfterPublish = true cfg.Pipeline.Workspace.CleanupAfterPublish = true failKey := seed.sessionPrefix + "current/run_id.txt" publishStageImpl, err := stage.Select("publish") if err != nil { t.Fatalf("Select(publish) error = %v", err) } _, err = executeStages(context.Background(), cfg, []stage.Stage{publishStageImpl}, RunOptions{ Env: &Env{ObjectStore: &failKeyStore{delegate: &storage.FakeBackend{}, failKey: failKey}}, }) if err == nil || !strings.Contains(err.Error(), "current run pointer") { t.Fatalf("executeStages() error = %v, want current-run-pointer failure", err) } assertExists(t, seed.spoolAudioDir) assertExists(t, seed.runWorkDir) } type cleanupSeed struct { runWorkDir string otherRunDir string spoolAudioDir string localSourceAudio string previousCachePath string sessionPrefix string } func cleanupFixtureConfig(t *testing.T) (*config.Config, cleanupSeed) { t.Helper() cfg := testConfig(t) cfg.Pipeline.Publish = &config.PublishConfig{Enabled: boolPtr(true), UploadRun: boolPtr(true)} cfg.Pipeline.Spool.Root = filepath.Join(t.TempDir(), "spool") runID := "20260516T010203Z-1a2b3c4d" runWorkDir := artifacts.SessionRunRootForCampaign(cfg.Pipeline.Workspace.Root, cfg.Session.Campaign, cfg.Session.SessionID, runID) otherRunDir := artifacts.SessionRunRootForCampaign(cfg.Pipeline.Workspace.Root, cfg.Session.Campaign, cfg.Session.SessionID, "20260516T010204Z-5e6f7a8b") spoolAudioDir := artifacts.SessionSpoolAudioDir(cfg.Pipeline.Spool.Root, cfg.Session.Campaign, cfg.Session.SessionID, runID) previousCachePath := artifacts.SessionPreviousArtifactPathForCampaign( cfg.Pipeline.Workspace.Root, cfg.Session.Campaign, cfg.Session.SessionID, "session_recap.md", ) mustWriteFile(t, filepath.Join(runWorkDir, "manifest.json"), "{}\n") mustWriteFile(t, filepath.Join(runWorkDir, "logs", "stage.log"), "log\n") mustWriteFile(t, filepath.Join(otherRunDir, "logs", "stage.log"), "other\n") mustWriteFile(t, filepath.Join(spoolAudioDir, "speaker.flac"), "flac\n") mustWriteFile(t, previousCachePath, "# previous recap\n") localSourceAudio := filepath.Join(filepath.Dir(cfg.SessionPath), "audio", "alice.flac") mustWriteFile(t, localSourceAudio, "source\n") seed := manifest.New(cfg.Session.SessionID, time.Now().UTC()) seed.Campaign = cfg.Session.Campaign seed.RunID = runID seed.LocalWorkDir = runWorkDir seed.LocalSpoolDir = spoolAudioDir seed.S3Bucket = "my-dnd-archive" seed.S3SessionPrefix = "dnd/campaigns/sample-campaign/sessions/2026-05-03/" seed.S3RunPrefix = seed.S3SessionPrefix + "runs/" + runID + "/" store := &manifest.LocalStore{} if err := os.MkdirAll(filepath.Dir(manifestPathFor(cfg)), 0o755); err != nil { t.Fatalf("MkdirAll() error = %v", err) } if err := store.Save(context.Background(), manifestPathFor(cfg), seed); err != nil { t.Fatalf("seed manifest save error = %v", err) } return cfg, cleanupSeed{ runWorkDir: runWorkDir, otherRunDir: otherRunDir, spoolAudioDir: spoolAudioDir, localSourceAudio: localSourceAudio, previousCachePath: previousCachePath, sessionPrefix: seed.S3SessionPrefix, } } func publishStageCleanupFixture(t *testing.T) (*config.Config, cleanupSeed, string) { t.Helper() cfg, seed := cleanupFixtureConfig(t) runID := "20260516T010203Z-1a2b3c4d" cfg.Pipeline.Storage.S3 = &config.StorageS3Config{ Bucket: "my-dnd-archive", RootPrefix: "dnd", } cfg.Pipeline.Publish = &config.PublishConfig{ Enabled: boolPtr(true), UploadRun: boolPtr(true), Outputs: []config.PublishOutputRule{ {Source: "narratio.transcript.final_trimmed", Dest: "transcripts/final.trimmed.json", Required: boolPtr(true)}, {Source: "narratio.artifact.session_recap", Dest: "artifacts/session_recap.md", Required: boolPtr(true)}, }, } cfg.Pipeline.Scriptorium = &config.ScriptoriumConfig{ Artifacts: map[string]config.ScriptoriumArtifactConfig{ "session_recap": { OutputPath: "artifacts/session_recap.md", }, }, } writePublishFixtureRunFiles( t, seed.runWorkDir, artifacts.SessionWorkDirForCampaign(cfg.Pipeline.Workspace.Root, cfg.Session.Campaign, cfg.Session.SessionID), ) store := &manifest.LocalStore{} seedManifest, err := store.Load(context.Background(), manifestPathFor(cfg)) if err != nil { t.Fatalf("Load() error = %v", err) } for _, name := range []string{"prepare", "transcribe", "merge", "polish", "normalize", "trim", "extract", "render", "analyze"} { seedManifest.MarkStageSucceeded(name, time.Now().UTC(), nil) } seedManifest.S3SessionPrefix = artifacts.S3SessionPrefix("dnd", cfg.Session.Campaign, cfg.Session.SessionID) seedManifest.S3RunPrefix = artifacts.S3RunPrefix(seedManifest.S3SessionPrefix, runID) if err := store.Save(context.Background(), manifestPathFor(cfg), seedManifest); err != nil { t.Fatalf("Save() error = %v", err) } seed.sessionPrefix = seedManifest.S3SessionPrefix return cfg, seed, runID } func writePublishFixtureRunFiles(t *testing.T, runWorkDir, sessionRoot string) { t.Helper() mustWriteFile(t, filepath.Join(runWorkDir, "prepare", "inputs", "session.yml"), "session_id: 2026-05-03\n") mustWriteFile(t, filepath.Join(runWorkDir, "transcribe", "outputs", "transcripts", "raw", "speaker.json"), "{}\n") mustWriteFile(t, filepath.Join(runWorkDir, "trim", "outputs", "transcripts", "final.trimmed.json"), "{\"segments\":[]}\n") mustWriteFile(t, filepath.Join(runWorkDir, "analyze", "outputs", "artifacts", "session_recap.md"), "# recap\n") mustWriteFile(t, filepath.Join(runWorkDir, "polish", "reports", "audita.report.json"), "{}\n") mustWriteFile(t, filepath.Join(runWorkDir, "merge", "config", "seriatim.generated.yml"), "key: value\n") mustWriteFile(t, filepath.Join(runWorkDir, "logs", "audita.stderr.log"), "stderr\n") mustWriteFile(t, filepath.Join(runWorkDir, "manifest.json"), "{}\n") mustWriteFile(t, filepath.Join(sessionRoot, "transcripts", "final.trimmed.json"), "{\"segments\":[]}\n") mustWriteFile(t, filepath.Join(sessionRoot, "artifacts", "session_recap.md"), "# recap\n") } type failKeyStore struct { delegate *storage.FakeBackend failKey string } func (s *failKeyStore) List(ctx context.Context, prefix string) ([]storage.ObjectInfo, error) { return s.delegate.List(ctx, prefix) } func (s *failKeyStore) Download(ctx context.Context, key, localPath string) error { return s.delegate.Download(ctx, key, localPath) } func (s *failKeyStore) Upload(ctx context.Context, localPath, key string, opts storage.UploadOptions) (storage.ObjectInfo, error) { if strings.TrimSpace(key) == strings.TrimSpace(s.failKey) { return storage.ObjectInfo{}, errors.New("forced upload failure") } return s.delegate.Upload(ctx, localPath, key, opts) } func (s *failKeyStore) Exists(ctx context.Context, key string) (bool, error) { return s.delegate.Exists(ctx, key) } func assertExists(t *testing.T, path string) { t.Helper() if _, err := os.Stat(path); err != nil { t.Fatalf("expected path to exist %q: %v", path, err) } } func assertMissing(t *testing.T, path string) { t.Helper() if _, err := os.Stat(path); !os.IsNotExist(err) { t.Fatalf("expected path to be removed %q, stat err=%v", path, err) } }