diff --git a/internal/app/runner.go b/internal/app/runner.go index 7bc5f59..644b0fd 100644 --- a/internal/app/runner.go +++ b/internal/app/runner.go @@ -8,9 +8,9 @@ import ( "path/filepath" "strings" - "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/scriptorium" "gitea.maximumdirect.net/eric/narratio/internal/adapters/seriatim" "gitea.maximumdirect.net/eric/narratio/internal/adapters/storage" "gitea.maximumdirect.net/eric/narratio/internal/adapters/whisperx" @@ -72,8 +72,8 @@ func executeStages(ctx context.Context, cfg *config.Config, stages []stage.Stage } env.Audita = runner } - if env.Analyzer == nil { - env.Analyzer = &analyzer.NoopRunner{} + if env.Scriptorium == nil { + env.Scriptorium = scriptorium.NewSubprocessRunner() } if env.Storage == nil { env.Storage = &storage.NoopBackend{} diff --git a/internal/app/runner_test.go b/internal/app/runner_test.go index 9ef3c03..b24e525 100644 --- a/internal/app/runner_test.go +++ b/internal/app/runner_test.go @@ -9,9 +9,9 @@ import ( "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/scriptorium" "gitea.maximumdirect.net/eric/narratio/internal/adapters/seriatim" "gitea.maximumdirect.net/eric/narratio/internal/adapters/storage" "gitea.maximumdirect.net/eric/narratio/internal/adapters/whisperx" @@ -114,6 +114,15 @@ func TestExecuteStagesPlaceholderSuccessUpdatesManifest(t *testing.T) { } continue } + if name == "analyze" { + if sr.Metadata == nil || sr.Metadata["stage"] != "analyze" { + t.Fatalf("analyze metadata missing stage=analyze: %#v", sr.Metadata) + } + if sr.Metadata["skipped"] != true { + t.Fatalf("analyze metadata missing skipped=true when scriptorium is unconfigured: %#v", sr.Metadata) + } + continue + } if sr.Metadata == nil || sr.Metadata["placeholder"] != true { t.Fatalf("stage %q missing placeholder metadata", name) } @@ -284,7 +293,7 @@ func TestAdapterBackedStageFailureMarksManifestFailed(t *testing.T) { {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: "analyze", env: &Env{Scriptorium: &scriptorium.FakeRunner{RunErr: 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")}}}, } @@ -347,6 +356,29 @@ func TestAdapterBackedStageFailureMarksManifestFailed(t *testing.T) { t.Fatalf("write glossary: %v", err) } } + if tc.name == "analyze" { + paths, ensureErr := artifactStore.EnsureLayout(cfg.Session.SessionID) + if ensureErr != nil { + t.Fatalf("EnsureLayout() error = %v", ensureErr) + } + if err := os.WriteFile(filepath.Join(paths.TranscriptsDir, "processed.json"), []byte(`{"segments":[]}`), 0o644); err != nil { + t.Fatalf("write processed transcript: %v", err) + } + cfg.Pipeline.Scriptorium = &config.ScriptoriumConfig{ + Binary: "scriptorium", + Timeout: "10m", + Artifacts: map[string]config.ScriptoriumArtifactConfig{ + "session_recap": { + Enabled: true, + PromptID: "dnd.session_recap", + OutputPath: "artifacts/session_recap.md", + Inputs: map[string]config.ScriptoriumInputConfig{ + "transcript": {Source: "processed_transcript", Required: true}, + }, + }, + }, + } + } _, runErr := executeStages(context.Background(), cfg, []stage.Stage{selected}, RunOptions{Env: tc.env}) if runErr == nil { diff --git a/internal/stage/analyze.go b/internal/stage/analyze.go new file mode 100644 index 0000000..c0d4ee9 --- /dev/null +++ b/internal/stage/analyze.go @@ -0,0 +1,462 @@ +package stage + +import ( + "context" + "fmt" + "os" + "path/filepath" + "sort" + "strings" + "time" + + "gitea.maximumdirect.net/eric/narratio/internal/adapters/scriptorium" + "gitea.maximumdirect.net/eric/narratio/internal/artifacts" + "gitea.maximumdirect.net/eric/narratio/internal/config" + "gitea.maximumdirect.net/eric/narratio/internal/manifest" +) + +type analyzeStage struct{} + +func (analyzeStage) Name() string { return "analyze" } + +func (analyzeStage) Declares() IODecl { + return IODecl{ + Inputs: []artifacts.Ref{ + {Kind: "transcript_processed", Category: "transcripts", RelativePath: "transcripts/processed.json"}, + {Kind: "artifact", Category: "artifacts", RelativePath: "artifacts/session_recap.md"}, + }, + Outputs: []artifacts.Ref{ + {Kind: "session_recap", Category: "artifacts", RelativePath: "artifacts/session_recap.md"}, + }, + } +} + +func (analyzeStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*StageResult, error) { + if env == nil || env.Config == nil { + return nil, fmt.Errorf("analyze: stage environment config is required") + } + if env.ArtifactStore == nil { + return nil, fmt.Errorf("analyze: artifact store is required") + } + if env.Config.Pipeline == nil || env.Config.Session == nil { + return nil, fmt.Errorf("analyze: resolved config must include pipeline and session") + } + if env.Scriptorium == nil { + return nil, fmt.Errorf("analyze: scriptorium adapter is required") + } + + var sessionID string + if m != nil { + sessionID = strings.TrimSpace(m.SessionID) + } + if sessionID == "" { + sessionID = strings.TrimSpace(env.Config.Session.SessionID) + } + if sessionID == "" { + return nil, fmt.Errorf("analyze: session id is required") + } + + paths := env.ArtifactStore.SessionPaths(sessionID) + if env.Config.Pipeline.Scriptorium == nil { + return &StageResult{ + Metadata: map[string]any{ + "stage": "analyze", + "skipped": true, + "reason": "pipeline.scriptorium is not configured", + }, + }, nil + } + + artifactName, artifactCfg, skipReason, err := selectAnalyzeArtifact(env.Config.Pipeline.Scriptorium) + if err != nil { + return nil, fmt.Errorf("analyze: %w", err) + } + if skipReason != "" { + return &StageResult{ + Metadata: map[string]any{ + "stage": "analyze", + "skipped": true, + "reason": skipReason, + }, + }, nil + } + + processedTranscriptPath, processedSource, err := discoverProcessedTranscript(m, paths) + if err != nil { + return nil, fmt.Errorf("analyze: resolve processed transcript: %w", err) + } + if processedTranscriptPath == "" { + return nil, fmt.Errorf("analyze: processed transcript input is required") + } + if err := validateProcessedTranscriptOutput(processedTranscriptPath); err != nil { + return nil, fmt.Errorf("analyze: processed transcript %q invalid: %w", processedTranscriptPath, err) + } + + inputPaths := map[string]string{} + omittedOptionalInputs := []string{} + sessionDir := filepath.Dir(strings.TrimSpace(env.Config.SessionPath)) + inputNames := sortedScriptoriumInputNames(artifactCfg.Inputs) + for _, inputName := range inputNames { + inputCfg := artifactCfg.Inputs[inputName] + resolvedPath, resolved, resolveErr := resolveScriptoriumInput(inputName, inputCfg, processedTranscriptPath, paths, sessionDir) + if resolveErr != nil { + return nil, fmt.Errorf("analyze: resolve input %q: %w", inputName, resolveErr) + } + if !resolved { + if inputCfg.Required { + return nil, fmt.Errorf("analyze: required input %q could not be resolved", inputName) + } + omittedOptionalInputs = append(omittedOptionalInputs, inputName) + continue + } + inputPaths[inputName] = resolvedPath + } + + vars, err := buildScriptoriumVars(artifactCfg.Vars, env.Config.Session) + if err != nil { + return nil, fmt.Errorf("analyze: resolve vars: %w", err) + } + + outputPath, err := resolveScriptoriumOutputPath(paths, artifactCfg.OutputPath) + if err != nil { + return nil, fmt.Errorf("analyze: resolve output path: %w", err) + } + stdoutLogPath := filepath.Join(paths.LogsDir, "scriptorium."+artifactName+".stdout.log") + stderrLogPath := filepath.Join(paths.LogsDir, "scriptorium."+artifactName+".stderr.log") + generatedConfigPath := filepath.Join(paths.ConfigDir, "scriptorium."+artifactName+".generated.yml") + + timeout, err := resolveScriptoriumTimeout(env.Config.Pipeline.Scriptorium.Timeout, artifactCfg.Timeout) + if err != nil { + return nil, fmt.Errorf("analyze: resolve timeout: %w", err) + } + + req := scriptorium.RunArtifactRequest{ + Binary: env.Config.Pipeline.Scriptorium.Binary, + ConfigPath: env.Config.Pipeline.Scriptorium.ConfigPath, + PromptID: artifactCfg.PromptID, + ProfileID: artifactCfg.ProfileID, + InputPaths: inputPaths, + Vars: vars, + OutputPath: outputPath, + StdoutLogPath: stdoutLogPath, + StderrLogPath: stderrLogPath, + GeneratedConfigPath: generatedConfigPath, + Timeout: timeout, + } + + res, runErr := env.Scriptorium.RunArtifact(ctx, req) + if runErr != nil { + if res.ValidationFailed { + return nil, fmt.Errorf( + "analyze: scriptorium validation failed (prompt_id=%q, output_path=%q, exit_code=%d, stdout_log=%q, stderr_log=%q): %w", + req.PromptID, + coalesceString(res.OutputPath, req.OutputPath), + res.ExitCode, + coalesceString(res.StdoutLogPath, req.StdoutLogPath), + coalesceString(res.StderrLogPath, req.StderrLogPath), + runErr, + ) + } + return nil, fmt.Errorf("analyze: scriptorium run failed: %w", runErr) + } + if res.ValidationFailed { + return nil, fmt.Errorf("analyze: scriptorium run returned validation_failed=true") + } + + finalOutputPath := coalesceString(res.OutputPath, req.OutputPath) + if err := requireNonEmptyFile(finalOutputPath, "session recap output"); err != nil { + return nil, fmt.Errorf("analyze: %w", err) + } + + artifactRef := artifacts.Ref{ + Kind: artifactName, + Category: "artifacts", + SessionID: sessionID, + AbsolutePath: finalOutputPath, + } + + meta := map[string]any{ + "stage": "analyze", + "artifact_name": artifactName, + "prompt_id": artifactCfg.PromptID, + "profile_id": artifactCfg.ProfileID, + "output_path": finalOutputPath, + "generated_config_path": generatedConfigPath, + "stdout_log_path": stdoutLogPath, + "stderr_log_path": stderrLogPath, + "input_paths": inputPaths, + "input_names": sortedMapKeys(inputPaths), + "omitted_optional_inputs": omittedOptionalInputs, + "vars": vars, + "timeout": timeout.String(), + "binary": env.Config.Pipeline.Scriptorium.Binary, + "config_path": env.Config.Pipeline.Scriptorium.ConfigPath, + "processed_transcript_path": processedTranscriptPath, + "processed_transcript_source": processedSource, + "adapter_exit_code": res.ExitCode, + "adapter_duration_ms": res.Duration.Milliseconds(), + "adapter_command_mode": res.CommandMode, + "adapter_prompt_id": res.PromptID, + "adapter_profile_id": res.ProfileID, + "adapter_validation_failed": res.ValidationFailed, + "adapter_output_path": res.OutputPath, + "adapter_generated_config": res.GeneratedConfigPath, + "adapter_stdout_log_path": res.StdoutLogPath, + "adapter_stderr_log_path": res.StderrLogPath, + } + if res.Metadata != nil { + meta["adapter_metadata"] = res.Metadata + } + + return &StageResult{ + Outputs: []artifacts.Ref{artifactRef}, + Logs: []string{stdoutLogPath, stderrLogPath}, + GeneratedConfigs: []string{generatedConfigPath}, + Metadata: meta, + }, nil +} + +func selectAnalyzeArtifact(cfg *config.ScriptoriumConfig) (string, config.ScriptoriumArtifactConfig, string, error) { + if cfg == nil { + return "", config.ScriptoriumArtifactConfig{}, "pipeline.scriptorium is not configured", nil + } + + enabled := []string{} + for name, artifact := range cfg.Artifacts { + if artifact.Enabled { + enabled = append(enabled, name) + } + } + sort.Strings(enabled) + if len(enabled) == 0 { + return "", config.ScriptoriumArtifactConfig{}, "no enabled scriptorium artifacts configured", nil + } + + sessionRecapCfg, ok := cfg.Artifacts["session_recap"] + if !ok || !sessionRecapCfg.Enabled { + return "", config.ScriptoriumArtifactConfig{}, "", fmt.Errorf("only artifacts.session_recap is supported in this analyze implementation; enabled=%s", strings.Join(enabled, ",")) + } + return "session_recap", sessionRecapCfg, "", nil +} + +func discoverProcessedTranscript(m *manifest.Manifest, paths artifacts.SessionPaths) (string, string, error) { + candidates := []string{} + if m != nil && m.Stages != nil { + if sr := m.Stages["polish"]; sr != nil { + for _, out := range sr.Outputs { + if out.Kind != "transcript_processed" { + continue + } + p := strings.TrimSpace(out.LocalPath) + if p == "" { + continue + } + resolved := artifacts.ResolveSessionLocalPathForRead(paths, p) + candidates = append(candidates, filepath.Clean(resolved)) + } + } + } + deduped := dedupeAndSortPaths(candidates) + for _, p := range deduped { + if info, err := os.Stat(p); err == nil && !info.IsDir() { + return p, "manifest.polish.outputs", nil + } + } + + fallback := filepath.Join(paths.TranscriptsDir, "processed.json") + if info, err := os.Stat(fallback); err == nil && !info.IsDir() { + return filepath.Clean(fallback), "fallback.transcripts_dir", nil + } + if len(deduped) > 0 { + return deduped[0], "manifest.polish.outputs", nil + } + return "", "", nil +} + +func resolveScriptoriumInput( + inputName string, + inputCfg config.ScriptoriumInputConfig, + processedTranscriptPath string, + paths artifacts.SessionPaths, + sessionDir string, +) (string, bool, error) { + switch strings.TrimSpace(inputCfg.Source) { + case "processed_transcript": + return processedTranscriptPath, true, nil + case "previous_session_artifact": + if strings.TrimSpace(inputCfg.Path) == "" { + return "", false, nil + } + resolved := resolveInputPathForRead(paths, sessionDir, inputCfg.Path) + if err := requireFile(resolved, "scriptorium input "+inputName); err != nil { + return "", false, nil + } + return resolved, true, nil + default: + return "", false, fmt.Errorf("unsupported source %q", inputCfg.Source) + } +} + +func resolveInputPathForRead(paths artifacts.SessionPaths, sessionDir, pathValue string) string { + trimmed := strings.TrimSpace(pathValue) + if trimmed == "" { + return "" + } + if filepath.IsAbs(trimmed) { + return filepath.Clean(trimmed) + } + candidates := []string{} + if strings.TrimSpace(sessionDir) != "" { + candidates = append(candidates, filepath.Clean(filepath.Join(sessionDir, trimmed))) + } + candidates = append(candidates, filepath.Clean(artifacts.ResolveSessionLocalPathForRead(paths, trimmed))) + for _, c := range candidates { + if info, err := os.Stat(c); err == nil && !info.IsDir() { + return c + } + } + return candidates[0] +} + +func resolveScriptoriumOutputPath(paths artifacts.SessionPaths, configured string) (string, error) { + outputPath := strings.TrimSpace(configured) + if outputPath == "" { + outputPath = "artifacts/session_recap.md" + } + if filepath.IsAbs(outputPath) { + return filepath.Clean(outputPath), nil + } + + rel := filepath.Clean(outputPath) + if rel == "." || rel == "" { + return "", fmt.Errorf("relative output path is required") + } + if rel == ".." || strings.HasPrefix(rel, ".."+string(filepath.Separator)) { + return "", fmt.Errorf("relative output path escapes session root: %q", outputPath) + } + return filepath.Join(paths.Root, rel), nil +} + +func resolveScriptoriumTimeout(topLevel, artifact string) (time.Duration, error) { + raw := strings.TrimSpace(artifact) + if raw == "" { + raw = strings.TrimSpace(topLevel) + } + if raw == "" { + raw = "10m" + } + d, err := time.ParseDuration(raw) + if err != nil { + return 0, fmt.Errorf("parse duration %q: %w", raw, err) + } + if d <= 0 { + return 0, fmt.Errorf("duration must be > 0") + } + return d, nil +} + +func buildScriptoriumVars(varsCfg map[string]any, session *config.SessionConfig) (map[string]string, error) { + if len(varsCfg) == 0 { + return nil, nil + } + vars := map[string]string{} + keys := make([]string, 0, len(varsCfg)) + for key := range varsCfg { + keys = append(keys, key) + } + sort.Strings(keys) + + for _, key := range keys { + value := varsCfg[key] + switch typed := value.(type) { + case bool: + if !typed { + continue + } + derived, ok, err := deriveSessionVarValue(key, session) + if err != nil { + return nil, err + } + if ok { + vars[key] = derived + } + case string: + vars[key] = typed + default: + return nil, fmt.Errorf("var %q has unsupported type %T", key, value) + } + } + if len(vars) == 0 { + return nil, nil + } + return vars, nil +} + +func deriveSessionVarValue(name string, session *config.SessionConfig) (string, bool, error) { + switch name { + case "session_id": + if session == nil || strings.TrimSpace(session.SessionID) == "" { + return "", false, nil + } + return strings.TrimSpace(session.SessionID), true, nil + case "session_date": + if session == nil || strings.TrimSpace(session.Date) == "" { + return "", false, nil + } + return strings.TrimSpace(session.Date), true, nil + case "campaign_name": + if session == nil || strings.TrimSpace(session.Campaign) == "" { + return "", false, nil + } + return strings.TrimSpace(session.Campaign), true, nil + case "previous_session_id": + return "", false, nil + default: + return "", false, fmt.Errorf("unsupported boolean var %q", name) + } +} + +func requireNonEmptyFile(path string, label string) error { + if err := requireFile(path, label); err != nil { + return err + } + info, err := os.Stat(path) + if err != nil { + return fmt.Errorf("%s %q stat failed: %w", label, path, err) + } + if info.Size() <= 0 { + return fmt.Errorf("%s %q is empty", label, path) + } + return nil +} + +func sortedScriptoriumInputNames(inputs map[string]config.ScriptoriumInputConfig) []string { + if len(inputs) == 0 { + return nil + } + names := make([]string, 0, len(inputs)) + for name := range inputs { + names = append(names, name) + } + sort.Strings(names) + return names +} + +func sortedMapKeys(values map[string]string) []string { + if len(values) == 0 { + return nil + } + keys := make([]string, 0, len(values)) + for key := range values { + keys = append(keys, key) + } + sort.Strings(keys) + return keys +} + +func coalesceString(primary, fallback string) string { + if strings.TrimSpace(primary) != "" { + return primary + } + return fallback +} diff --git a/internal/stage/analyze_test.go b/internal/stage/analyze_test.go new file mode 100644 index 0000000..7986b20 --- /dev/null +++ b/internal/stage/analyze_test.go @@ -0,0 +1,379 @@ +package stage + +import ( + "context" + "errors" + "os" + "path/filepath" + "strings" + "testing" + "time" + + "gitea.maximumdirect.net/eric/narratio/internal/adapters/scriptorium" + "gitea.maximumdirect.net/eric/narratio/internal/artifacts" + "gitea.maximumdirect.net/eric/narratio/internal/config" + "gitea.maximumdirect.net/eric/narratio/internal/manifest" +) + +func TestAnalyzeGeneratesSessionRecapFromProcessedTranscript(t *testing.T) { + env, m, fake := setupAnalyzeEnv(t) + paths := env.ArtifactStore.SessionPaths(m.SessionID) + writeAnalyzeFile(t, filepath.Join(paths.TranscriptsDir, "processed.json"), `{"segments":[]}`) + + result, err := (analyzeStage{}).Run(context.Background(), env, m) + if err != nil { + t.Fatalf("Run() error = %v", err) + } + if len(fake.RunRequests) != 1 { + t.Fatalf("scriptorium run requests = %d, want 1", len(fake.RunRequests)) + } + + req := fake.RunRequests[0] + if req.PromptID != "dnd.session_recap" { + t.Fatalf("prompt id = %q, want dnd.session_recap", req.PromptID) + } + if req.ProfileID != "local-quality" { + t.Fatalf("profile id = %q, want local-quality", req.ProfileID) + } + if req.InputPaths["transcript"] != filepath.Join(paths.TranscriptsDir, "processed.json") { + t.Fatalf("transcript input = %q, want processed transcript path", req.InputPaths["transcript"]) + } + if req.OutputPath != filepath.Join(paths.ArtifactsDir, "session_recap.md") { + t.Fatalf("output path = %q, want %q", req.OutputPath, filepath.Join(paths.ArtifactsDir, "session_recap.md")) + } + if req.StdoutLogPath != filepath.Join(paths.LogsDir, "scriptorium.session_recap.stdout.log") { + t.Fatalf("stdout log path = %q", req.StdoutLogPath) + } + if req.StderrLogPath != filepath.Join(paths.LogsDir, "scriptorium.session_recap.stderr.log") { + t.Fatalf("stderr log path = %q", req.StderrLogPath) + } + if req.GeneratedConfigPath != filepath.Join(paths.ConfigDir, "scriptorium.session_recap.generated.yml") { + t.Fatalf("generated config path = %q", req.GeneratedConfigPath) + } + if req.Timeout != 2*time.Minute { + t.Fatalf("timeout = %s, want 2m", req.Timeout) + } + + if len(result.Outputs) != 1 || result.Outputs[0].Kind != "session_recap" { + t.Fatalf("outputs = %#v, want one session_recap output", result.Outputs) + } + if len(result.Logs) != 2 { + t.Fatalf("logs = %#v, want stdout+stderr logs", result.Logs) + } + if len(result.GeneratedConfigs) != 1 { + t.Fatalf("generated configs = %#v, want one path", result.GeneratedConfigs) + } + if result.Metadata["stage"] != "analyze" { + t.Fatalf("metadata stage = %#v, want analyze", result.Metadata["stage"]) + } +} + +func TestAnalyzeOmitsOptionalPreviousRecapWhenUnavailable(t *testing.T) { + env, m, fake := setupAnalyzeEnv(t) + paths := env.ArtifactStore.SessionPaths(m.SessionID) + writeAnalyzeFile(t, filepath.Join(paths.TranscriptsDir, "processed.json"), `{"segments":[]}`) + + env.Config.Pipeline.Scriptorium.Artifacts["session_recap"] = config.ScriptoriumArtifactConfig{ + Enabled: true, + PromptID: "dnd.session_recap", + ProfileID: "local-quality", + OutputPath: "artifacts/session_recap.md", + Timeout: "2m", + Inputs: map[string]config.ScriptoriumInputConfig{ + "transcript": { + Source: "processed_transcript", + Required: true, + }, + "previous_recap": { + Source: "previous_session_artifact", + Artifact: "session_recap", + Path: "", + Required: false, + }, + }, + Vars: map[string]any{ + "session_id": true, + }, + } + + _, err := (analyzeStage{}).Run(context.Background(), env, m) + if err != nil { + t.Fatalf("Run() error = %v", err) + } + req := fake.RunRequests[0] + if _, exists := req.InputPaths["previous_recap"]; exists { + t.Fatalf("previous_recap input should be omitted when optional+unavailable, got %q", req.InputPaths["previous_recap"]) + } +} + +func TestAnalyzeIncludesPreviousRecapWhenConfiguredAndAvailable(t *testing.T) { + env, m, fake := setupAnalyzeEnv(t) + paths := env.ArtifactStore.SessionPaths(m.SessionID) + writeAnalyzeFile(t, filepath.Join(paths.TranscriptsDir, "processed.json"), `{"segments":[]}`) + + previousRecapPath := filepath.Join(filepath.Dir(env.Config.SessionPath), "previous", "session_recap.md") + writeAnalyzeFile(t, previousRecapPath, "previous recap\n") + + env.Config.Pipeline.Scriptorium.Artifacts["session_recap"] = config.ScriptoriumArtifactConfig{ + Enabled: true, + PromptID: "dnd.session_recap", + ProfileID: "local-quality", + OutputPath: "artifacts/session_recap.md", + Timeout: "2m", + Inputs: map[string]config.ScriptoriumInputConfig{ + "transcript": { + Source: "processed_transcript", + Required: true, + }, + "previous_recap": { + Source: "previous_session_artifact", + Artifact: "session_recap", + Path: "./previous/session_recap.md", + Required: false, + }, + }, + Vars: map[string]any{ + "session_id": true, + }, + } + + _, err := (analyzeStage{}).Run(context.Background(), env, m) + if err != nil { + t.Fatalf("Run() error = %v", err) + } + req := fake.RunRequests[0] + if req.InputPaths["previous_recap"] != previousRecapPath { + t.Fatalf("previous_recap input = %q, want %q", req.InputPaths["previous_recap"], previousRecapPath) + } +} + +func TestAnalyzeFailsWhenRequiredPreviousRecapMissing(t *testing.T) { + env, m, _ := setupAnalyzeEnv(t) + paths := env.ArtifactStore.SessionPaths(m.SessionID) + writeAnalyzeFile(t, filepath.Join(paths.TranscriptsDir, "processed.json"), `{"segments":[]}`) + + env.Config.Pipeline.Scriptorium.Artifacts["session_recap"] = config.ScriptoriumArtifactConfig{ + Enabled: true, + PromptID: "dnd.session_recap", + OutputPath: "artifacts/session_recap.md", + Inputs: map[string]config.ScriptoriumInputConfig{ + "transcript": { + Source: "processed_transcript", + Required: true, + }, + "previous_recap": { + Source: "previous_session_artifact", + Path: "./missing/previous_recap.md", + Required: true, + }, + }, + } + + _, err := (analyzeStage{}).Run(context.Background(), env, m) + if err == nil { + t.Fatal("expected error, got nil") + } + if !strings.Contains(err.Error(), `required input "previous_recap"`) { + t.Fatalf("error = %q, want required input context", err.Error()) + } +} + +func TestAnalyzeFailsWhenProcessedTranscriptMissing(t *testing.T) { + env, m, _ := setupAnalyzeEnv(t) + _, err := (analyzeStage{}).Run(context.Background(), env, m) + if err == nil { + t.Fatal("expected error, got nil") + } + if !strings.Contains(err.Error(), "processed transcript input is required") { + t.Fatalf("error = %q, want missing processed transcript context", err.Error()) + } +} + +func TestAnalyzeFailsWhenProcessedTranscriptInvalidJSON(t *testing.T) { + env, m, _ := setupAnalyzeEnv(t) + paths := env.ArtifactStore.SessionPaths(m.SessionID) + writeAnalyzeFile(t, filepath.Join(paths.TranscriptsDir, "processed.json"), `{not-json`) + + _, err := (analyzeStage{}).Run(context.Background(), env, m) + if err == nil { + t.Fatal("expected error, got nil") + } + if !strings.Contains(err.Error(), "decode json") { + t.Fatalf("error = %q, want decode json context", err.Error()) + } +} + +func TestAnalyzeFailsWhenProcessedTranscriptMissingSegmentsArray(t *testing.T) { + env, m, _ := setupAnalyzeEnv(t) + paths := env.ArtifactStore.SessionPaths(m.SessionID) + writeAnalyzeFile(t, filepath.Join(paths.TranscriptsDir, "processed.json"), `{"not_segments":[]}`) + + _, err := (analyzeStage{}).Run(context.Background(), env, m) + if err == nil { + t.Fatal("expected error, got nil") + } + if !strings.Contains(err.Error(), "top-level segments is required") { + t.Fatalf("error = %q, want segments guidance", err.Error()) + } +} + +func TestAnalyzeRecordsRefsAndMetadata(t *testing.T) { + env, m, _ := setupAnalyzeEnv(t) + paths := env.ArtifactStore.SessionPaths(m.SessionID) + writeAnalyzeFile(t, filepath.Join(paths.TranscriptsDir, "processed.json"), `{"segments":[]}`) + + result, err := (analyzeStage{}).Run(context.Background(), env, m) + if err != nil { + t.Fatalf("Run() error = %v", err) + } + if len(result.Outputs) != 1 { + t.Fatalf("outputs len = %d, want 1", len(result.Outputs)) + } + if result.Outputs[0].AbsolutePath != filepath.Join(paths.ArtifactsDir, "session_recap.md") { + t.Fatalf("output path = %q, want session recap path", result.Outputs[0].AbsolutePath) + } + if len(result.Logs) != 2 { + t.Fatalf("logs = %#v, want two logs", result.Logs) + } + if len(result.GeneratedConfigs) != 1 { + t.Fatalf("generated configs = %#v, want one", result.GeneratedConfigs) + } + if result.Metadata["adapter_command_mode"] != scriptorium.CommandModeRun { + t.Fatalf("adapter command mode = %#v, want run", result.Metadata["adapter_command_mode"]) + } + if result.Metadata["prompt_id"] != "dnd.session_recap" { + t.Fatalf("prompt_id = %#v, want dnd.session_recap", result.Metadata["prompt_id"]) + } +} + +func TestAnalyzeHandlesAdapterError(t *testing.T) { + env, m, fake := setupAnalyzeEnv(t) + paths := env.ArtifactStore.SessionPaths(m.SessionID) + writeAnalyzeFile(t, filepath.Join(paths.TranscriptsDir, "processed.json"), `{"segments":[]}`) + fake.RunErr = errors.New("adapter boom") + + _, err := (analyzeStage{}).Run(context.Background(), env, m) + if err == nil { + t.Fatal("expected error, got nil") + } + if !strings.Contains(err.Error(), "scriptorium run failed") { + t.Fatalf("error = %q, want adapter failure context", err.Error()) + } +} + +func TestAnalyzeHandlesValidationFailedResultAsError(t *testing.T) { + env, m, fake := setupAnalyzeEnv(t) + paths := env.ArtifactStore.SessionPaths(m.SessionID) + writeAnalyzeFile(t, filepath.Join(paths.TranscriptsDir, "processed.json"), `{"segments":[]}`) + fake.RunResult = scriptorium.ArtifactResult{ + ValidationFailed: true, + ExitCode: 2, + CommandMode: scriptorium.CommandModeRun, + } + + _, err := (analyzeStage{}).Run(context.Background(), env, m) + if err == nil { + t.Fatal("expected error, got nil") + } + if !strings.Contains(err.Error(), "validation_failed=true") { + t.Fatalf("error = %q, want validation_failed context", err.Error()) + } +} + +func TestAnalyzeSkipsWhenNoEnabledScriptoriumArtifactsConfigured(t *testing.T) { + env, m, _ := setupAnalyzeEnv(t) + paths := env.ArtifactStore.SessionPaths(m.SessionID) + writeAnalyzeFile(t, filepath.Join(paths.TranscriptsDir, "processed.json"), `{"segments":[]}`) + + env.Config.Pipeline.Scriptorium.Artifacts["session_recap"] = config.ScriptoriumArtifactConfig{ + Enabled: false, + } + + result, err := (analyzeStage{}).Run(context.Background(), env, m) + if err != nil { + t.Fatalf("Run() error = %v", err) + } + if result.Metadata["skipped"] != true { + t.Fatalf("metadata = %#v, want skipped=true", result.Metadata) + } +} + +func setupAnalyzeEnv(t *testing.T) (*Env, *manifest.Manifest, *scriptorium.FakeRunner) { + t.Helper() + workspace := t.TempDir() + cfgDir := t.TempDir() + sessionPath := filepath.Join(cfgDir, "session.yml") + pipelinePath := filepath.Join(cfgDir, "pipeline.yml") + + writeAnalyzeFile(t, pipelinePath, "workspace:\n root: "+workspace+"\n") + writeAnalyzeFile(t, sessionPath, "session_id: 2026-05-03\n") + + fake := &scriptorium.FakeRunner{} + cfg := &config.Config{ + PipelinePath: pipelinePath, + SessionPath: sessionPath, + Pipeline: &config.PipelineConfig{ + Workspace: config.WorkspaceConfig{Root: workspace}, + Scriptorium: &config.ScriptoriumConfig{ + Binary: "scriptorium", + ConfigPath: "/etc/scriptorium/config.yml", + Timeout: "10m", + Artifacts: map[string]config.ScriptoriumArtifactConfig{ + "session_recap": { + Enabled: true, + PromptID: "dnd.session_recap", + ProfileID: "local-quality", + OutputPath: "artifacts/session_recap.md", + Timeout: "2m", + Inputs: map[string]config.ScriptoriumInputConfig{ + "transcript": { + Source: "processed_transcript", + Required: true, + }, + "previous_recap": { + Source: "previous_session_artifact", + Artifact: "session_recap", + Path: "", + Required: false, + }, + }, + Vars: map[string]any{ + "session_id": true, + "session_date": true, + "campaign_name": true, + "output_kind": "session_recap", + }, + }, + }, + }, + }, + Session: &config.SessionConfig{ + SessionID: "2026-05-03", + Campaign: "Icewind Dale", + Date: "2026-05-03", + }, + } + + store := artifacts.NewLocalStore(workspace) + if _, err := store.EnsureLayout("2026-05-03"); err != nil { + t.Fatalf("EnsureLayout() error = %v", err) + } + + env := &Env{ + Config: cfg, + ArtifactStore: store, + Scriptorium: fake, + } + m := manifest.New("2026-05-03", time.Now().UTC()) + return env, m, fake +} + +func writeAnalyzeFile(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) + } +} diff --git a/internal/stage/placeholders.go b/internal/stage/placeholders.go index 16b35b2..77fac75 100644 --- a/internal/stage/placeholders.go +++ b/internal/stage/placeholders.go @@ -5,7 +5,6 @@ import ( "fmt" "path/filepath" - "gitea.maximumdirect.net/eric/narratio/internal/adapters/analyzer" "gitea.maximumdirect.net/eric/narratio/internal/adapters/notify" "gitea.maximumdirect.net/eric/narratio/internal/adapters/storage" "gitea.maximumdirect.net/eric/narratio/internal/artifacts" @@ -51,25 +50,6 @@ func (s placeholderStage) Run(ctx context.Context, env *Env, m *manifest.Manifes paths := env.ArtifactStore.SessionPaths(sessionID) switch s.name { - case "analyze": - if env.Analyzer != nil { - req := analyzer.AnalyzeRequest{ - ArtifactType: "session-log", - ProcessedTranscriptPath: filepath.Join(paths.TranscriptsDir, "processed.json"), - ContextReferences: []string{"previous-session"}, - OutputPath: filepath.Join(paths.ArtifactsDir, "session-log.md"), - GeneratedConfigPath: filepath.Join(paths.ConfigDir, "analyzer.session-log.generated.yml"), - StdoutLogPath: filepath.Join(paths.LogsDir, "analyzer.session-log.stdout.log"), - StderrLogPath: filepath.Join(paths.LogsDir, "analyzer.session-log.stderr.log"), - } - resp, err := env.Analyzer.Run(ctx, req) - if err != nil { - return nil, fmt.Errorf("placeholder analyze adapter call failed: %w", err) - } - result.Outputs = append(result.Outputs, artifacts.Ref{Kind: "artifact", Category: "artifacts", SessionID: sessionID, AbsolutePath: resp.ArtifactPath}) - result.Logs = append(result.Logs, req.StdoutLogPath, req.StderrLogPath) - result.GeneratedConfigs = append(result.GeneratedConfigs, req.GeneratedConfigPath) - } case "archive": if env.Storage != nil { req := storage.ArchiveRequest{ @@ -113,7 +93,7 @@ func All() []Stage { placeholderStage{name: "normalize"}, mergeStage{}, polishStage{}, - placeholderStage{name: "analyze"}, + analyzeStage{}, placeholderStage{name: "archive"}, placeholderStage{name: "notify"}, } diff --git a/internal/stage/placeholders_test.go b/internal/stage/placeholders_test.go index 8f0bc3d..f3d789c 100644 --- a/internal/stage/placeholders_test.go +++ b/internal/stage/placeholders_test.go @@ -9,9 +9,9 @@ import ( "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/scriptorium" "gitea.maximumdirect.net/eric/narratio/internal/adapters/seriatim" "gitea.maximumdirect.net/eric/narratio/internal/adapters/storage" "gitea.maximumdirect.net/eric/narratio/internal/adapters/whisperx" @@ -41,7 +41,7 @@ func TestStagesReturnExpectedMetadata(t *testing.T) { wf := &whisperx.FakeClient{} sf := &seriatim.FakeRunner{} af := &audita.FakeRunner{} - anz := &analyzer.FakeRunner{} + sc := &scriptorium.FakeRunner{} st := &storage.FakeBackend{} nf := ¬ify.FakeSender{} @@ -64,7 +64,7 @@ func TestStagesReturnExpectedMetadata(t *testing.T) { WhisperX: wf, Seriatim: sf, Audita: af, - Analyzer: anz, + Scriptorium: sc, Storage: st, Notifier: nf, } @@ -111,6 +111,15 @@ func TestStagesReturnExpectedMetadata(t *testing.T) { } continue } + if s.Name() == "analyze" { + if result.Metadata["stage"] != "analyze" { + t.Fatalf("analyze metadata = %#v, want stage=analyze", result.Metadata) + } + if result.Metadata["skipped"] != true { + t.Fatalf("analyze metadata = %#v, want skipped=true when scriptorium is unconfigured", result.Metadata) + } + continue + } if result.Metadata["placeholder"] != true { t.Fatalf("stage %q missing placeholder metadata", s.Name()) } @@ -125,8 +134,8 @@ func TestStagesReturnExpectedMetadata(t *testing.T) { if len(af.Requests) != 1 { t.Fatalf("audita calls = %d, want 1", len(af.Requests)) } - if len(anz.Requests) != 1 { - t.Fatalf("analyzer calls = %d, want 1", len(anz.Requests)) + if len(sc.RunRequests) != 0 { + t.Fatalf("scriptorium run calls = %d, want 0 when scriptorium config is absent", len(sc.RunRequests)) } if len(st.Requests) != 1 { t.Fatalf("storage calls = %d, want 1", len(st.Requests)) @@ -142,7 +151,6 @@ func TestPlaceholderAdapterErrorPropagation(t *testing.T) { env *Env wantErr string }{ - {stageName: "analyze", env: &Env{Analyzer: &analyzer.FakeRunner{Err: errors.New("anerr")}}, wantErr: "analyze"}, {stageName: "archive", env: &Env{Storage: &storage.FakeBackend{Err: errors.New("sterr")}}, wantErr: "archive"}, {stageName: "notify", env: &Env{Notifier: ¬ify.FakeSender{Err: errors.New("nerr")}}, wantErr: "notify"}, } diff --git a/internal/stage/stage.go b/internal/stage/stage.go index 7a6bbc0..6ba4861 100644 --- a/internal/stage/stage.go +++ b/internal/stage/stage.go @@ -7,6 +7,7 @@ import ( "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/scriptorium" "gitea.maximumdirect.net/eric/narratio/internal/adapters/seriatim" "gitea.maximumdirect.net/eric/narratio/internal/adapters/storage" "gitea.maximumdirect.net/eric/narratio/internal/adapters/whisperx" @@ -22,12 +23,13 @@ type Env struct { ManifestStore manifest.Store Logger *slog.Logger - WhisperX whisperx.Client - Seriatim seriatim.Runner - Audita audita.Runner - Analyzer analyzer.Runner - Storage storage.Backend - Notifier notify.Sender + WhisperX whisperx.Client + Seriatim seriatim.Runner + Audita audita.Runner + Scriptorium scriptorium.Runner + Analyzer analyzer.Runner + Storage storage.Backend + Notifier notify.Sender } // IODecl declares the intended input/output artifact kinds for a stage.