diff --git a/internal/stage/merge.go b/internal/stage/merge.go index 4e591ff..4cfcfff 100644 --- a/internal/stage/merge.go +++ b/internal/stage/merge.go @@ -206,6 +206,10 @@ func normalizeMergeInputs(ctx context.Context, env *Env, rawInputs []string, pat logs := make([]string, 0, len(rawInputs)*2) configs := make([]string, 0, len(rawInputs)) meta := make([]normalizeMergeInputMeta, 0, len(rawInputs)) + normalizedDir := filepath.Join(paths.TranscriptsRawDir, "normalized") + if err := os.MkdirAll(normalizedDir, 0o755); err != nil { + return nil, nil, nil, nil, fmt.Errorf("merge: ensure normalized transcripts directory %q: %w", normalizedDir, err) + } var timeout time.Duration timeoutRaw := strings.TrimSpace(env.Config.Pipeline.Seriatim.Timeout) @@ -218,7 +222,7 @@ func normalizeMergeInputs(ctx context.Context, env *Env, rawInputs []string, pat } for _, input := range rawInputs { base := strings.TrimSuffix(filepath.Base(input), filepath.Ext(input)) - outPath := filepath.Join(paths.TranscriptsRawDir, "normalized", base+".normalized.json") + outPath := filepath.Join(normalizedDir, base+".normalized.json") stdoutPath := filepath.Join(paths.LogsDir, "seriatim.normalize."+base+".stdout.log") stderrPath := filepath.Join(paths.LogsDir, "seriatim.normalize."+base+".stderr.log") cfgPath := filepath.Join(paths.ConfigDir, "seriatim.normalize."+base+".generated.yml") diff --git a/internal/stage/merge_test.go b/internal/stage/merge_test.go index 26d0e97..0877121 100644 --- a/internal/stage/merge_test.go +++ b/internal/stage/merge_test.go @@ -2,6 +2,7 @@ package stage import ( "context" + "fmt" "os" "path/filepath" "strings" @@ -14,6 +15,22 @@ import ( "gitea.maximumdirect.net/eric/narratio/internal/manifest" ) +type normalizeDirAssertingRunner struct { + *seriatim.FakeRunner + ExpectedDir string +} + +func (r *normalizeDirAssertingRunner) Normalize(ctx context.Context, req seriatim.NormalizeRequest) (seriatim.NormalizeResult, error) { + info, err := os.Stat(r.ExpectedDir) + if err != nil { + return seriatim.NormalizeResult{}, fmt.Errorf("normalized directory check failed for %q: %w", r.ExpectedDir, err) + } + if !info.IsDir() { + return seriatim.NormalizeResult{}, fmt.Errorf("normalized path %q exists but is not a directory", r.ExpectedDir) + } + return r.FakeRunner.Normalize(ctx, req) +} + func TestMergeStageMergesRawTranscriptsAndRecordsMetadata(t *testing.T) { env, m := setupMergeEnv(t) paths := env.ArtifactStore.SessionPaths(m.SessionID) @@ -205,6 +222,30 @@ func TestMergeStageResolvesWorkspaceQualifiedManifestOutputsWithoutDuplication(t } } +func TestMergeStageCreatesNormalizedRawDirectoryBeforeNormalize(t *testing.T) { + env, m := setupMergeEnv(t) + paths := env.ArtifactStore.SessionPaths(m.SessionID) + + rawPath := filepath.Join(paths.TranscriptsRawDir, "alice.json") + writeFile(t, rawPath, `{"segments":[]}`) + writeFile(t, filepath.Join(paths.InputsDir, "speakers.yml"), "match: []\n") + writeFile(t, filepath.Join(paths.InputsDir, "autocorrect.yml"), "rules: []\n") + + normalizedDir := filepath.Join(paths.TranscriptsRawDir, "normalized") + if err := os.RemoveAll(normalizedDir); err != nil { + t.Fatalf("remove normalized dir: %v", err) + } + + env.Seriatim = &normalizeDirAssertingRunner{ + FakeRunner: &seriatim.FakeRunner{}, + ExpectedDir: normalizedDir, + } + + if _, err := (mergeStage{}).Run(context.Background(), env, m); err != nil { + t.Fatalf("merge.Run() error = %v", err) + } +} + func TestMergeStageFailsWhenNormalizeAdapterFails(t *testing.T) { env, m := setupMergeEnv(t) paths := env.ArtifactStore.SessionPaths(m.SessionID)