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, campaignPath, sessionPath := writeValidConfigFilesWithScriptoriumArtifacts(t, workspaceRoot) fake := &storage.FakeBackend{} cfg, sessionPrefix, manifestKey, runIDKey := seedRestoreCommittedState(t, fake, pipelinePath, campaignPath, 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{ "session", "restore", cfg.Session.SessionID, "--config", pipelinePath, "--campaign-file", campaignPath, "--session", sessionPath, }, &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", "analyze", cfg.Session.SessionID, "--config", pipelinePath, "--campaign-file", campaignPath, "--session", sessionPath, "--force", "--artifacts", "player_handout", }, &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 TestRestoreThenAnalyzeUsesRestoredPreviousCacheWithoutObjectStore(t *testing.T) { workspaceRoot := t.TempDir() pipelinePath, campaignPath, sessionPath := writeValidConfigFiles(t, workspaceRoot) appendRestoreWorkflowScriptoriumConfig(t, pipelinePath, ` scriptorium: binary: scriptorium artifacts: session_recap: enabled: true prompt_id: dnd.session_recap output_path: artifacts/session_recap.md inputs: transcript: source: narratio.transcript.final_trimmed required: true previous_recap: source: narratio.previous_session.artifact.session_recap required: true `) appendRestoreWorkflowScriptoriumConfig(t, sessionPath, ` previous_session_id: 2026-04-26 `) fakeStore := &storage.FakeBackend{} cfg, sessionPrefix, manifestKey, runIDKey := seedRestoreCommittedState(t, fakeStore, pipelinePath, campaignPath, sessionPath) seedRestoreObject(fakeStore, runIDKey, []byte("20260519T010203Z-a1b2c3d4\n")) seedRestoreObject(fakeStore, manifestKey, restoreWorkflowManifestJSON(t, cfg.Session.SessionID, cfg.Session.Campaign)) seedRestoreObject(fakeStore, sessionPrefix+"transcripts/final.trimmed.json", []byte(`{"segments":[]}`+"\n")) seedRestorePreviousCurrent(t, fakeStore, cfg, "# previous recap\n") restoreWithStoreAndRealPhases(t, fakeStore) var stdout bytes.Buffer var stderr bytes.Buffer restoreCode := Execute( []string{ "session", "restore", cfg.Session.SessionID, "--config", pipelinePath, "--campaign-file", campaignPath, "--session", sessionPath, }, &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", "final.trimmed.json"), `{"segments":[]}`+"\n") previousManifestBytes, err := os.ReadFile(filepath.Join(sessionRoot, "previous", "manifest.json")) if err != nil { t.Fatalf("read restored previous manifest: %v", err) } if !strings.Contains(string(previousManifestBytes), `"session_id":"2026-04-26"`) { t.Fatalf("restored previous manifest = %q, want previous session id", string(previousManifestBytes)) } mustReadEquals(t, filepath.Join(sessionRoot, "previous", "artifacts", "session_recap.md"), "# previous recap\n") scriptoriumFake := &scriptorium.FakeRunner{} origExecuteStagesFn := executeStagesFn origObjectStoreFn := newObjectStoreFromConfigFn objectStoreConstructed := false t.Cleanup(func() { executeStagesFn = origExecuteStagesFn newObjectStoreFromConfigFn = origObjectStoreFn }) 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 = scriptoriumFake return executeStages(ctx, cfg, stages, opts) } newObjectStoreFromConfigFn = func(context.Context, *config.Config) (storage.ObjectStore, error) { objectStoreConstructed = true return nil, context.Canceled } stdout.Reset() stderr.Reset() runStageCode := Execute( []string{ "run-stage", "analyze", cfg.Session.SessionID, "--config", pipelinePath, "--campaign-file", campaignPath, "--session", sessionPath, "--force", "--artifacts", "session_recap", }, &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 objectStoreConstructed { t.Fatal("analyze run-stage should not construct object store for previous-session input resolution") } if len(scriptoriumFake.RunRequests) != 1 { t.Fatalf("scriptorium run requests = %d, want 1", len(scriptoriumFake.RunRequests)) } req := scriptoriumFake.RunRequests[0] if got := req.InputPaths["transcript"]; got != filepath.Join(sessionRoot, "transcripts", "final.trimmed.json") { t.Fatalf("transcript input = %q, want trimmed transcript path", got) } if got := req.InputPaths["previous_recap"]; got != filepath.Join(sessionRoot, "previous", "artifacts", "session_recap.md") { t.Fatalf("previous_recap input = %q, want restored previous cache path", got) } } 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 } func appendRestoreWorkflowScriptoriumConfig(t *testing.T, pipelinePath, extra string) { t.Helper() f, err := os.OpenFile(pipelinePath, os.O_APPEND|os.O_WRONLY, 0) if err != nil { t.Fatalf("open pipeline config for append: %v", err) } defer f.Close() if _, err := f.WriteString(extra); err != nil { t.Fatalf("append pipeline config: %v", err) } }