422 lines
16 KiB
Go
422 lines
16 KiB
Go
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 := mustPreviousArtifactPathForCampaign(t,
|
|
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)
|
|
}
|
|
}
|