diff --git a/internal/app/restore_workflow_test.go b/internal/app/restore_workflow_test.go new file mode 100644 index 0000000..15f001e --- /dev/null +++ b/internal/app/restore_workflow_test.go @@ -0,0 +1,178 @@ +package app + +import ( + "bytes" + "context" + "os" + "path/filepath" + "strings" + "testing" + "time" + + "gitea.maximumdirect.net/eric/narratio/internal/adapters/scriptorium" + "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" +) + +func TestRestoreThenRunStageForceAnalyzeUsesRestoredDurableState(t *testing.T) { + workspaceRoot := t.TempDir() + pipelinePath, sessionPath := writeValidConfigFilesWithScriptoriumArtifacts(t, workspaceRoot) + + fake := &storage.FakeBackend{} + cfg, sessionPrefix, manifestKey, runIDKey := seedRestoreCommittedState(t, fake, pipelinePath, sessionPath) + seedRestoreObject(fake, runIDKey, []byte("20260519T010203Z-a1b2c3d4\n")) + seedRestoreObject(fake, manifestKey, restoreWorkflowManifestJSON(t, cfg.Session.SessionID, cfg.Session.Campaign)) + seedRestoreObject(fake, sessionPrefix+"transcripts/full.json", []byte(`{"segments":[1,2,3]}`+"\n")) + seedRestoreObject(fake, sessionPrefix+"artifacts/session_recap.md", []byte("# restored recap\n")) + + restoreWithStoreAndRealPhases(t, fake) + + origExecuteStagesFn := executeStagesFn + t.Cleanup(func() { + executeStagesFn = origExecuteStagesFn + }) + executeStagesFn = func(ctx context.Context, cfg *config.Config, stages []stage.Stage, opts RunOptions) (*RunSummary, error) { + if opts.Env == nil { + opts.Env = &Env{} + } + opts.Env.Scriptorium = &scriptorium.NoopRunner{} + return executeStages(ctx, cfg, stages, opts) + } + + var stdout bytes.Buffer + var stderr bytes.Buffer + restoreCode := Execute( + []string{ + "restore", + "--config", pipelinePath, + "--session", sessionPath, + "--session-id", cfg.Session.SessionID, + }, + &stdout, + &stderr, + ) + if restoreCode != 0 { + t.Fatalf("restore exit code = %d, want 0; stderr=%q", restoreCode, stderr.String()) + } + if stderr.Len() != 0 { + t.Fatalf("restore stderr = %q, want empty", stderr.String()) + } + + sessionRoot := artifacts.SessionWorkDirForCampaign(workspaceRoot, cfg.Session.Campaign, cfg.Session.SessionID) + mustReadEquals(t, filepath.Join(sessionRoot, "transcripts", "full.json"), `{"segments":[1,2,3]}`+"\n") + mustReadEquals(t, filepath.Join(sessionRoot, "artifacts", "session_recap.md"), "# restored recap\n") + + manifestStore := &manifest.LocalStore{} + sessionManifestPath := artifacts.SessionManifestPathForCampaign(workspaceRoot, cfg.Session.Campaign, cfg.Session.SessionID) + beforeAnalyze, err := manifestStore.Load(context.Background(), sessionManifestPath) + if err != nil { + t.Fatalf("load restored session manifest: %v", err) + } + upstreamCompletedAt := map[string]time.Time{} + for _, stageName := range []string{"prepare", "transcribe", "merge", "polish", "normalize", "trim"} { + rec := beforeAnalyze.Stages[stageName] + if rec == nil || rec.Status != manifest.StatusSucceeded || rec.CompletedAt == nil { + t.Fatalf("restored manifest stage %q = %#v, want succeeded with completion timestamp", stageName, rec) + } + upstreamCompletedAt[stageName] = *rec.CompletedAt + } + + stdout.Reset() + stderr.Reset() + runStageCode := Execute( + []string{ + "run-stage", + "--config", pipelinePath, + "--session", sessionPath, + "--session-id", cfg.Session.SessionID, + "--force", + "--artifacts", "player_handout", + "analyze", + }, + &stdout, + &stderr, + ) + if runStageCode != 0 { + t.Fatalf("run-stage exit code = %d, want 0; stderr=%q", runStageCode, stderr.String()) + } + if stderr.Len() != 0 { + t.Fatalf("run-stage stderr = %q, want empty", stderr.String()) + } + if !strings.Contains(stdout.String(), "stage=analyze executed=1 skipped=0 force=true") { + t.Fatalf("run-stage stdout = %q, want analyze execution summary", stdout.String()) + } + + playerHandoutPath := filepath.Join(sessionRoot, "artifacts", "player_handout.md") + if _, err := os.Stat(playerHandoutPath); err != nil { + t.Fatalf("restored analyze output %q missing: %v", playerHandoutPath, err) + } + + afterAnalyze, err := manifestStore.Load(context.Background(), sessionManifestPath) + if err != nil { + t.Fatalf("load session manifest after run-stage analyze: %v", err) + } + for _, stageName := range []string{"prepare", "transcribe", "merge", "polish", "normalize", "trim"} { + rec := afterAnalyze.Stages[stageName] + if rec == nil || rec.Status != manifest.StatusSucceeded || rec.CompletedAt == nil { + t.Fatalf("post-analyze manifest stage %q = %#v, want succeeded with completion timestamp", stageName, rec) + } + if !rec.CompletedAt.Equal(upstreamCompletedAt[stageName]) { + t.Fatalf( + "stage %q completion changed: before=%s after=%s", + stageName, + upstreamCompletedAt[stageName].Format(time.RFC3339Nano), + rec.CompletedAt.Format(time.RFC3339Nano), + ) + } + } + analyzeRec := afterAnalyze.Stages["analyze"] + if analyzeRec == nil || analyzeRec.Status != manifest.StatusSucceeded { + t.Fatalf("post-analyze stage record = %#v, want succeeded", analyzeRec) + } + + runManifestPaths, err := filepath.Glob(filepath.Join(sessionRoot, "runs", "*", "manifest.json")) + if err != nil { + t.Fatalf("glob run manifests: %v", err) + } + if len(runManifestPaths) != 1 { + t.Fatalf("run manifest count = %d, want 1; paths=%v", len(runManifestPaths), runManifestPaths) + } + runManifest, err := manifestStore.LoadRun(context.Background(), runManifestPaths[0]) + if err != nil { + t.Fatalf("load run manifest %q: %v", runManifestPaths[0], err) + } + if len(runManifest.RequestedStages) != 1 || runManifest.RequestedStages[0] != "analyze" { + t.Fatalf("run manifest requested_stages = %#v, want [analyze]", runManifest.RequestedStages) + } + if runManifest.Stages["analyze"] == nil || runManifest.Stages["analyze"].Status != manifest.StatusSucceeded { + t.Fatalf("run manifest analyze stage = %#v, want succeeded", runManifest.Stages["analyze"]) + } + if runManifest.Stages["prepare"] != nil { + t.Fatalf("run manifest should not include upstream prepare stage, got %#v", runManifest.Stages["prepare"]) + } +} + +func restoreWorkflowManifestJSON(t *testing.T, sessionID, campaign string) []byte { + t.Helper() + store := &manifest.LocalStore{} + now := time.Date(2026, 5, 19, 23, 0, 0, 0, time.UTC) + m := manifest.New(sessionID, now) + m.Campaign = campaign + m.RunID = "20260519T010203Z-a1b2c3d4" + stages := []string{"prepare", "transcribe", "merge", "polish", "normalize", "trim"} + for i, stageName := range stages { + m.MarkStageSucceeded(stageName, now.Add(time.Duration(i+1)*time.Minute), nil) + } + path := filepath.Join(t.TempDir(), "manifest.json") + if err := store.Save(context.Background(), path, m); err != nil { + t.Fatalf("save workflow manifest fixture: %v", err) + } + data, err := os.ReadFile(path) + if err != nil { + t.Fatalf("read workflow manifest fixture: %v", err) + } + return data +} diff --git a/internal/app/resume.go b/internal/app/resume.go index fec720d..3306bfc 100644 --- a/internal/app/resume.go +++ b/internal/app/resume.go @@ -76,7 +76,7 @@ func Resume(ctx context.Context, args []string, out io.Writer) error { } } - summary, err := executeStages(ctx, cfg, selected, RunOptions{ + summary, err := executeStagesFn(ctx, cfg, selected, RunOptions{ Force: force, SelectedArtifacts: normalizedArtifacts, }) diff --git a/internal/app/run.go b/internal/app/run.go index 9401556..703bdcc 100644 --- a/internal/app/run.go +++ b/internal/app/run.go @@ -58,7 +58,7 @@ func Run(ctx context.Context, args []string, out io.Writer) error { } stages := BuildFullPlan() - summary, err := executeStages(ctx, cfg, stages, RunOptions{ + summary, err := executeStagesFn(ctx, cfg, stages, RunOptions{ Force: force, SelectedArtifacts: normalizedArtifacts, }) diff --git a/internal/app/run_stage.go b/internal/app/run_stage.go index 32ccc7f..2ff7cdd 100644 --- a/internal/app/run_stage.go +++ b/internal/app/run_stage.go @@ -66,7 +66,7 @@ func RunStage(ctx context.Context, args []string, out io.Writer) error { return fmt.Errorf("run-stage: %w", err) } - summary, err := executeStages(ctx, cfg, stages, RunOptions{ + summary, err := executeStagesFn(ctx, cfg, stages, RunOptions{ Force: force, SelectedArtifacts: normalizedArtifacts, }) diff --git a/internal/app/runner.go b/internal/app/runner.go index d12d7d5..77c35d3 100644 --- a/internal/app/runner.go +++ b/internal/app/runner.go @@ -36,6 +36,8 @@ type RunSummary struct { Skipped []string } +var executeStagesFn = executeStages + func executeStages(ctx context.Context, cfg *config.Config, stages []stage.Stage, opts RunOptions) (*RunSummary, error) { env := opts.Env if env == nil {