411 lines
14 KiB
Go
411 lines
14 KiB
Go
package app
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"os"
|
|
"path/filepath"
|
|
"strings"
|
|
"testing"
|
|
"time"
|
|
|
|
"gitea.maximumdirect.net/eric/narratio/internal/adapters/analyzer"
|
|
"gitea.maximumdirect.net/eric/narratio/internal/adapters/audita"
|
|
"gitea.maximumdirect.net/eric/narratio/internal/adapters/notify"
|
|
"gitea.maximumdirect.net/eric/narratio/internal/adapters/seriatim"
|
|
"gitea.maximumdirect.net/eric/narratio/internal/adapters/storage"
|
|
"gitea.maximumdirect.net/eric/narratio/internal/adapters/whisperx"
|
|
"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 failingStage struct {
|
|
name string
|
|
err error
|
|
}
|
|
|
|
func (s failingStage) Name() string { return s.name }
|
|
func (s failingStage) Declares() stage.IODecl { return stage.IODecl{} }
|
|
func (s failingStage) Run(_ context.Context, _ *stage.Env, _ *manifest.Manifest) (*stage.StageResult, error) {
|
|
return nil, s.err
|
|
}
|
|
|
|
type countingStage struct {
|
|
name string
|
|
runs *int
|
|
}
|
|
|
|
func (s countingStage) Name() string { return s.name }
|
|
func (s countingStage) Declares() stage.IODecl { return stage.IODecl{} }
|
|
func (s countingStage) Run(_ context.Context, _ *stage.Env, _ *manifest.Manifest) (*stage.StageResult, error) {
|
|
*s.runs = *s.runs + 1
|
|
return &stage.StageResult{Metadata: map[string]any{"counting": true}}, nil
|
|
}
|
|
|
|
func TestExecuteStagesPlaceholderSuccessUpdatesManifest(t *testing.T) {
|
|
cfg := testConfig(t)
|
|
|
|
summary, err := executeStages(context.Background(), cfg, BuildFullPlan(), RunOptions{})
|
|
if err != nil {
|
|
t.Fatalf("executeStages() error = %v", err)
|
|
}
|
|
if len(summary.StageNames) != 8 || len(summary.Executed) != 8 || len(summary.Skipped) != 0 {
|
|
t.Fatalf("summary = %#v, want all 8 executed", summary)
|
|
}
|
|
|
|
store := &manifest.LocalStore{}
|
|
m, err := store.Load(context.Background(), summary.ManifestPath)
|
|
if err != nil {
|
|
t.Fatalf("Load manifest error = %v", err)
|
|
}
|
|
|
|
for _, name := range []string{"prepare", "transcribe", "normalize", "merge", "polish", "analyze", "archive", "notify"} {
|
|
sr := m.Stages[name]
|
|
if sr == nil {
|
|
t.Fatalf("missing stage record %q", name)
|
|
}
|
|
if sr.Status != manifest.StatusSucceeded {
|
|
t.Fatalf("stage %q status = %q, want %q", name, sr.Status, manifest.StatusSucceeded)
|
|
}
|
|
if name == "prepare" {
|
|
if sr.Metadata == nil || sr.Metadata["prepared"] != true {
|
|
t.Fatalf("prepare metadata missing prepared=true: %#v", sr.Metadata)
|
|
}
|
|
continue
|
|
}
|
|
if name == "transcribe" {
|
|
if sr.Metadata == nil || sr.Metadata["stage"] != "transcribe" {
|
|
t.Fatalf("transcribe metadata missing stage=transcribe: %#v", sr.Metadata)
|
|
}
|
|
if len(sr.Outputs) == 0 {
|
|
t.Fatalf("transcribe outputs missing")
|
|
}
|
|
continue
|
|
}
|
|
if name == "merge" {
|
|
if sr.Metadata == nil || sr.Metadata["stage"] != "merge" {
|
|
t.Fatalf("merge metadata missing stage=merge: %#v", sr.Metadata)
|
|
}
|
|
if len(sr.Outputs) == 0 {
|
|
t.Fatalf("merge outputs missing")
|
|
}
|
|
if len(sr.Logs) == 0 {
|
|
t.Fatalf("merge logs missing")
|
|
}
|
|
if len(sr.GeneratedConfigs) == 0 {
|
|
t.Fatalf("merge generated configs missing")
|
|
}
|
|
continue
|
|
}
|
|
if name == "polish" {
|
|
if sr.Metadata == nil || sr.Metadata["stage"] != "polish" {
|
|
t.Fatalf("polish metadata missing stage=polish: %#v", sr.Metadata)
|
|
}
|
|
if len(sr.Outputs) == 0 {
|
|
t.Fatalf("polish outputs missing")
|
|
}
|
|
if len(sr.Logs) == 0 {
|
|
t.Fatalf("polish logs missing")
|
|
}
|
|
if len(sr.GeneratedConfigs) == 0 {
|
|
t.Fatalf("polish generated configs missing")
|
|
}
|
|
continue
|
|
}
|
|
if sr.Metadata == nil || sr.Metadata["placeholder"] != true {
|
|
t.Fatalf("stage %q missing placeholder metadata", name)
|
|
}
|
|
}
|
|
if len(m.Inputs) == 0 {
|
|
t.Fatalf("manifest inputs should be recorded by prepare")
|
|
}
|
|
|
|
if _, err := os.Stat(summary.ManifestPath); err != nil {
|
|
t.Fatalf("manifest file missing at %q: %v", summary.ManifestPath, err)
|
|
}
|
|
}
|
|
|
|
func TestExecuteStagesSkipSucceededWhenNotForced(t *testing.T) {
|
|
cfg := testConfig(t)
|
|
manifestPath := manifestPathFor(cfg)
|
|
store := &manifest.LocalStore{}
|
|
|
|
existing := manifest.New(cfg.Session.SessionID, time.Date(2026, 5, 3, 1, 0, 0, 0, time.UTC))
|
|
existing.MarkStageSucceeded("transcribe", time.Date(2026, 5, 3, 1, 1, 0, 0, time.UTC), []manifest.ArtifactRecord{{Kind: "transcript_raw", LocalPath: "existing.json"}})
|
|
existing.Stages["transcribe"].Metadata = map[string]any{"kept": true}
|
|
if err := os.MkdirAll(filepath.Dir(manifestPath), 0o755); err != nil {
|
|
t.Fatalf("MkdirAll() error = %v", err)
|
|
}
|
|
if err := store.Save(context.Background(), manifestPath, existing); err != nil {
|
|
t.Fatalf("Save manifest error = %v", err)
|
|
}
|
|
|
|
runs := 0
|
|
stageToRun := countingStage{name: "transcribe", runs: &runs}
|
|
summary, err := executeStages(context.Background(), cfg, []stage.Stage{stageToRun}, RunOptions{})
|
|
if err != nil {
|
|
t.Fatalf("executeStages() error = %v", err)
|
|
}
|
|
if runs != 0 {
|
|
t.Fatalf("runs = %d, want 0 due to skip", runs)
|
|
}
|
|
if len(summary.Executed) != 0 || len(summary.Skipped) != 1 || summary.Skipped[0] != "transcribe" {
|
|
t.Fatalf("summary = %#v, want skipped transcribe", summary)
|
|
}
|
|
|
|
loaded, err := store.Load(context.Background(), manifestPath)
|
|
if err != nil {
|
|
t.Fatalf("Load manifest error = %v", err)
|
|
}
|
|
sr := loaded.Stages["transcribe"]
|
|
if sr == nil || sr.Status != manifest.StatusSucceeded {
|
|
t.Fatalf("transcribe status = %#v, want succeeded", sr)
|
|
}
|
|
if sr.Metadata == nil || sr.Metadata["kept"] != true {
|
|
t.Fatalf("transcribe metadata = %#v, want preserved", sr.Metadata)
|
|
}
|
|
}
|
|
|
|
func TestExecuteStagesForceRerunsSucceeded(t *testing.T) {
|
|
cfg := testConfig(t)
|
|
manifestPath := manifestPathFor(cfg)
|
|
store := &manifest.LocalStore{}
|
|
|
|
existing := manifest.New(cfg.Session.SessionID, time.Date(2026, 5, 3, 1, 0, 0, 0, time.UTC))
|
|
existing.MarkStageSucceeded("transcribe", time.Date(2026, 5, 3, 1, 1, 0, 0, time.UTC), nil)
|
|
if err := os.MkdirAll(filepath.Dir(manifestPath), 0o755); err != nil {
|
|
t.Fatalf("MkdirAll() error = %v", err)
|
|
}
|
|
if err := store.Save(context.Background(), manifestPath, existing); err != nil {
|
|
t.Fatalf("Save manifest error = %v", err)
|
|
}
|
|
|
|
runs := 0
|
|
stageToRun := countingStage{name: "transcribe", runs: &runs}
|
|
summary, err := executeStages(context.Background(), cfg, []stage.Stage{stageToRun}, RunOptions{Force: true})
|
|
if err != nil {
|
|
t.Fatalf("executeStages() error = %v", err)
|
|
}
|
|
if runs != 1 {
|
|
t.Fatalf("runs = %d, want 1 with force", runs)
|
|
}
|
|
if len(summary.Executed) != 1 || summary.Executed[0] != "transcribe" || len(summary.Skipped) != 0 {
|
|
t.Fatalf("summary = %#v, want executed transcribe", summary)
|
|
}
|
|
}
|
|
|
|
func TestExecuteStagesFailureUpdatesManifest(t *testing.T) {
|
|
cfg := testConfig(t)
|
|
|
|
stages := []stage.Stage{
|
|
BuildFullPlan()[0],
|
|
failingStage{name: "transcribe", err: errors.New("boom")},
|
|
BuildFullPlan()[2],
|
|
}
|
|
|
|
summary, err := executeStages(context.Background(), cfg, stages, RunOptions{})
|
|
if err == nil {
|
|
t.Fatal("expected error, got nil")
|
|
}
|
|
if summary != nil {
|
|
t.Fatalf("summary = %#v, want nil on failure", summary)
|
|
}
|
|
if !strings.Contains(err.Error(), "stage \"transcribe\" failed") {
|
|
t.Fatalf("error = %q, want stage failure", err.Error())
|
|
}
|
|
|
|
manifestPath := manifestPathFor(cfg)
|
|
store := &manifest.LocalStore{}
|
|
m, loadErr := store.Load(context.Background(), manifestPath)
|
|
if loadErr != nil {
|
|
t.Fatalf("Load manifest error = %v", loadErr)
|
|
}
|
|
|
|
if got := m.Stages["prepare"]; got == nil || got.Status != manifest.StatusSucceeded {
|
|
t.Fatalf("prepare status = %#v, want succeeded", got)
|
|
}
|
|
if got := m.Stages["transcribe"]; got == nil || got.Status != manifest.StatusFailed {
|
|
t.Fatalf("transcribe status = %#v, want failed", got)
|
|
}
|
|
if got := m.Stages["normalize"]; got != nil {
|
|
t.Fatalf("normalize should not run, got %#v", got)
|
|
}
|
|
}
|
|
|
|
func TestExecuteStagesLoadsExistingManifest(t *testing.T) {
|
|
cfg := testConfig(t)
|
|
manifestPath := manifestPathFor(cfg)
|
|
store := &manifest.LocalStore{}
|
|
|
|
existing := manifest.New(cfg.Session.SessionID, time.Date(2026, 5, 3, 1, 0, 0, 0, time.UTC))
|
|
existing.MarkStageSucceeded("prepare", time.Date(2026, 5, 3, 1, 1, 0, 0, time.UTC), nil)
|
|
audioPath := filepath.Join(cfg.Pipeline.Workspace.Root, "work", cfg.Session.SessionID, "audio", "alice.flac")
|
|
if err := os.MkdirAll(filepath.Dir(audioPath), 0o755); err != nil {
|
|
t.Fatalf("MkdirAll() error = %v", err)
|
|
}
|
|
if err := os.WriteFile(audioPath, []byte("audio"), 0o644); err != nil {
|
|
t.Fatalf("WriteFile() error = %v", err)
|
|
}
|
|
existing.Inputs = append(existing.Inputs, manifest.InputRecord{
|
|
Kind: "audio",
|
|
Path: audioPath,
|
|
})
|
|
if err := os.MkdirAll(filepath.Dir(manifestPath), 0o755); err != nil {
|
|
t.Fatalf("MkdirAll() error = %v", err)
|
|
}
|
|
if err := store.Save(context.Background(), manifestPath, existing); err != nil {
|
|
t.Fatalf("Save manifest error = %v", err)
|
|
}
|
|
|
|
_, err := executeStages(context.Background(), cfg, []stage.Stage{BuildFullPlan()[1]}, RunOptions{})
|
|
if err != nil {
|
|
t.Fatalf("executeStages() error = %v", err)
|
|
}
|
|
|
|
loaded, err := store.Load(context.Background(), manifestPath)
|
|
if err != nil {
|
|
t.Fatalf("Load manifest error = %v", err)
|
|
}
|
|
if loaded.Stages["prepare"] == nil || loaded.Stages["prepare"].Status != manifest.StatusSucceeded {
|
|
t.Fatalf("existing stage prepare should remain succeeded")
|
|
}
|
|
if loaded.Stages["transcribe"] == nil || loaded.Stages["transcribe"].Status != manifest.StatusSucceeded {
|
|
t.Fatalf("transcribe should be succeeded after run")
|
|
}
|
|
}
|
|
|
|
func TestAdapterBackedStageFailureMarksManifestFailed(t *testing.T) {
|
|
cases := []struct {
|
|
name string
|
|
env *Env
|
|
}{
|
|
{name: "transcribe", env: &Env{WhisperX: &whisperx.FakeClient{Err: errors.New("transcribe fail")}}},
|
|
{name: "merge", env: &Env{Seriatim: &seriatim.FakeRunner{Err: errors.New("merge fail")}}},
|
|
{name: "polish", env: &Env{Audita: &audita.FakeRunner{Err: errors.New("polish fail")}}},
|
|
{name: "analyze", env: &Env{Analyzer: &analyzer.FakeRunner{Err: errors.New("analyze fail")}}},
|
|
{name: "archive", env: &Env{Storage: &storage.FakeBackend{Err: errors.New("archive fail")}}},
|
|
{name: "notify", env: &Env{Notifier: ¬ify.FakeSender{Err: errors.New("notify fail")}}},
|
|
}
|
|
|
|
for _, tc := range cases {
|
|
t.Run(tc.name, func(t *testing.T) {
|
|
cfg := testConfig(t)
|
|
selected, err := stage.Select(tc.name)
|
|
if err != nil {
|
|
t.Fatalf("Select() error = %v", err)
|
|
}
|
|
|
|
artifactStore := artifacts.NewLocalStore(cfg.Pipeline.Workspace.Root)
|
|
tc.env.Config = cfg
|
|
tc.env.ArtifactStore = artifactStore
|
|
tc.env.ManifestStore = &manifest.LocalStore{}
|
|
if tc.name == "transcribe" {
|
|
paths, ensureErr := artifactStore.EnsureLayout(cfg.Session.SessionID)
|
|
if ensureErr != nil {
|
|
t.Fatalf("EnsureLayout() error = %v", ensureErr)
|
|
}
|
|
audioPath := filepath.Join(paths.AudioDir, "alice.flac")
|
|
if err := os.WriteFile(audioPath, []byte("audio"), 0o644); err != nil {
|
|
t.Fatalf("write transcribe fixture audio: %v", err)
|
|
}
|
|
m := manifest.New(cfg.Session.SessionID, time.Now().UTC())
|
|
m.Inputs = append(m.Inputs, manifest.InputRecord{Kind: "audio", Path: audioPath})
|
|
if err := tc.env.ManifestStore.Save(context.Background(), manifestPathFor(cfg), m); err != nil {
|
|
t.Fatalf("seed manifest: %v", err)
|
|
}
|
|
}
|
|
if tc.name == "merge" {
|
|
paths, ensureErr := artifactStore.EnsureLayout(cfg.Session.SessionID)
|
|
if ensureErr != nil {
|
|
t.Fatalf("EnsureLayout() error = %v", ensureErr)
|
|
}
|
|
rawPath := filepath.Join(paths.TranscriptsRawDir, "alice.json")
|
|
if err := os.MkdirAll(filepath.Dir(rawPath), 0o755); err != nil {
|
|
t.Fatalf("mkdir raw dir: %v", err)
|
|
}
|
|
if err := os.WriteFile(rawPath, []byte(`{"segments":[]}`), 0o644); err != nil {
|
|
t.Fatalf("write raw transcript: %v", err)
|
|
}
|
|
if err := os.WriteFile(filepath.Join(paths.InputsDir, "speakers.yml"), []byte("match: []\n"), 0o644); err != nil {
|
|
t.Fatalf("write speakers: %v", err)
|
|
}
|
|
if err := os.WriteFile(filepath.Join(paths.InputsDir, "autocorrect.yml"), []byte("rules: []\n"), 0o644); err != nil {
|
|
t.Fatalf("write autocorrect: %v", err)
|
|
}
|
|
}
|
|
if tc.name == "polish" {
|
|
paths, ensureErr := artifactStore.EnsureLayout(cfg.Session.SessionID)
|
|
if ensureErr != nil {
|
|
t.Fatalf("EnsureLayout() error = %v", ensureErr)
|
|
}
|
|
if err := os.WriteFile(filepath.Join(paths.TranscriptsDir, "merged.json"), []byte(`{"segments":[]}`), 0o644); err != nil {
|
|
t.Fatalf("write merged transcript: %v", err)
|
|
}
|
|
if err := os.WriteFile(filepath.Join(paths.InputsDir, "glossary.yml"), []byte("terms: []\n"), 0o644); err != nil {
|
|
t.Fatalf("write glossary: %v", err)
|
|
}
|
|
}
|
|
|
|
_, runErr := executeStages(context.Background(), cfg, []stage.Stage{selected}, RunOptions{Env: tc.env})
|
|
if runErr == nil {
|
|
t.Fatal("expected error, got nil")
|
|
}
|
|
|
|
m, loadErr := tc.env.ManifestStore.Load(context.Background(), manifestPathFor(cfg))
|
|
if loadErr != nil {
|
|
t.Fatalf("load manifest error = %v", loadErr)
|
|
}
|
|
sr := m.Stages[tc.name]
|
|
if sr == nil {
|
|
t.Fatalf("missing stage record %q", tc.name)
|
|
}
|
|
if sr.Status != manifest.StatusFailed {
|
|
t.Fatalf("status = %q, want %q", sr.Status, manifest.StatusFailed)
|
|
}
|
|
})
|
|
}
|
|
}
|
|
|
|
func testConfig(t *testing.T) *config.Config {
|
|
t.Helper()
|
|
|
|
workspace := t.TempDir()
|
|
cfgDir := t.TempDir()
|
|
sessionPath := filepath.Join(cfgDir, "session.yml")
|
|
pipelinePath := filepath.Join(cfgDir, "pipeline.yml")
|
|
|
|
mustWriteFile(t, pipelinePath, "workspace:\n root: "+workspace+"\n")
|
|
mustWriteFile(t, sessionPath, "session_id: 2026-05-03\n")
|
|
mustWriteFile(t, filepath.Join(cfgDir, "speakers.yml"), "alice: alice.flac\n")
|
|
mustWriteFile(t, filepath.Join(cfgDir, "autocorrect.yml"), "[]\n")
|
|
mustWriteFile(t, filepath.Join(cfgDir, "glossary.yml"), "[]\n")
|
|
mustWriteFile(t, filepath.Join(cfgDir, "audio", "alice.flac"), "audio")
|
|
|
|
return &config.Config{
|
|
Pipeline: &config.PipelineConfig{Workspace: config.WorkspaceConfig{Root: workspace}},
|
|
PipelinePath: pipelinePath,
|
|
SessionPath: sessionPath,
|
|
Session: &config.SessionConfig{
|
|
SessionID: "2026-05-03",
|
|
Inputs: config.SessionInputsConfig{
|
|
AudioDir: "./audio",
|
|
SpeakersFile: "./speakers.yml",
|
|
AutocorrectFile: "./autocorrect.yml",
|
|
GlossaryFile: "./glossary.yml",
|
|
},
|
|
},
|
|
}
|
|
}
|
|
|
|
func mustWriteFile(t *testing.T, path, contents string) {
|
|
t.Helper()
|
|
if err := os.MkdirAll(filepath.Dir(path), 0o755); err != nil {
|
|
t.Fatalf("MkdirAll(%q): %v", path, err)
|
|
}
|
|
if err := os.WriteFile(path, []byte(contents), 0o644); err != nil {
|
|
t.Fatalf("WriteFile(%q): %v", path, err)
|
|
}
|
|
}
|