package app import ( "context" "errors" "os" "path/filepath" "reflect" "strings" "testing" "time" "gitea.maximumdirect.net/eric/narratio/internal/artifactmodel" "gitea.maximumdirect.net/eric/narratio/internal/artifacts" "gitea.maximumdirect.net/eric/narratio/internal/manifest" "gitea.maximumdirect.net/eric/narratio/internal/stage" ) type projectionStage struct { name string run func(*stage.Env, *manifest.Manifest) (*stage.StageResult, error) } func (s projectionStage) Name() string { return s.name } func (s projectionStage) Run(_ context.Context, env *stage.Env, m *manifest.Manifest) (*stage.StageResult, error) { return s.run(env, m) } func TestExecuteStagesProjectsSeparateSessionAndInvocationAnalyzeState(t *testing.T) { cfg := testConfig(t) oldAt := time.Date(2026, 5, 3, 9, 0, 0, 0, time.UTC) oldRecord := appAnalyzeRecord("player_handout", manifest.AnalyzeArtifactCurrent, "old-run", oldAt) stageToRun := projectionStage{name: "analyze", run: func(_ *stage.Env, m *manifest.Manifest) (*stage.StageResult, error) { newRecord := appAnalyzeRecord("session_recap", manifest.AnalyzeArtifactCurrent, m.RunID, time.Now().UTC()) staleRecord := appAnalyzeRecord("quest_log", manifest.AnalyzeArtifactStale, m.RunID, time.Now().UTC()) session := map[string]manifest.AnalyzeArtifactRecord{ "player_handout": oldRecord, "quest_log": staleRecord, "session_recap": newRecord, } return &stage.StageResult{ Logs: []string{"aggregate-analyze.log"}, AnalyzeState: &stage.AnalyzeStateProjection{ Session: session, Invocation: map[string]manifest.AnalyzeArtifactRecord{ "player_handout": oldRecord, "session_recap": newRecord, }, }, }, nil }} summary, err := executeStages(context.Background(), cfg, []stage.Stage{stageToRun}, RunOptions{}) if err != nil { t.Fatalf("executeStages() error = %v", err) } store := &manifest.LocalStore{} sessionManifest, err := store.Load(context.Background(), summary.ManifestPath) if err != nil { t.Fatalf("Load(session) error = %v", err) } analyze := sessionManifest.Stages["analyze"] if analyze.AnalyzeStateVersion != manifest.AnalyzeStateContractVersion || len(analyze.AnalyzeArtifacts) != 3 { t.Fatalf("session analyze state = %#v", analyze) } if got := analyzeArtifactOutputKeys(analyze.Outputs); !reflect.DeepEqual(got, []string{"player_handout", "session_recap"}) { t.Fatalf("session aggregate outputs = %#v, want current records only", got) } if len(analyze.Logs) != 1 || analyze.Logs[0] != "aggregate-analyze.log" { t.Fatalf("session aggregate logs = %#v", analyze.Logs) } runManifest, err := store.LoadRun(context.Background(), summary.RunManifestPath) if err != nil { t.Fatalf("LoadRun() error = %v", err) } runAnalyze := runManifest.Stages["analyze"] if len(runAnalyze.AnalyzeArtifacts) != 2 || runAnalyze.AnalyzeArtifacts["session_recap"].Status != manifest.AnalyzeArtifactCurrent { t.Fatalf("invocation analyze state = %#v", runAnalyze.AnalyzeArtifacts) } if got := analyzeArtifactOutputKeys(runAnalyze.Outputs); !reflect.DeepEqual(got, []string{"session_recap"}) { t.Fatalf("invocation outputs = %#v, want produced artifact only", got) } } func TestExecuteStagesPersistsRestrictedAnalyzeStateOnPartialError(t *testing.T) { cfg := testConfig(t) now := time.Date(2026, 5, 3, 9, 0, 0, 0, time.UTC) seed := manifest.New(cfg.Session.SessionID, now) seed.Campaign = cfg.Session.Campaign seed.MarkStageSucceeded("analyze", now, nil) seed.Stages["analyze"].AnalyzeStateVersion = manifest.AnalyzeStateContractVersion seed.Stages["analyze"].AnalyzeArtifacts = map[string]manifest.AnalyzeArtifactRecord{ "player_handout": appAnalyzeRecord("player_handout", manifest.AnalyzeArtifactCurrent, "old-run", now), } seed.MarkStageSucceeded("publish", now, nil) saveBoundedManifest(t, cfg, seed) stageToRun := projectionStage{name: "analyze", run: func(_ *stage.Env, m *manifest.Manifest) (*stage.StageResult, error) { unrelated := seed.Stages["analyze"].AnalyzeArtifacts["player_handout"] completed := appAnalyzeRecord("session_recap", manifest.AnalyzeArtifactCurrent, m.RunID, time.Now().UTC()) failed := appAnalyzeRecord("quest_log", manifest.AnalyzeArtifactFailed, m.RunID, time.Now().UTC()) session := map[string]manifest.AnalyzeArtifactRecord{ "player_handout": unrelated, "quest_log": failed, "session_recap": completed, } return &stage.StageResult{AnalyzeState: &stage.AnalyzeStateProjection{ Session: session, Invocation: map[string]manifest.AnalyzeArtifactRecord{ "quest_log": failed, "session_recap": completed, }, }}, errors.New("quest log failed") }} _, err := executeStages(context.Background(), cfg, []stage.Stage{stageToRun}, RunOptions{Force: true}) if err == nil || !strings.Contains(err.Error(), "quest log failed") { t.Fatalf("executeStages() error = %v", err) } store := &manifest.LocalStore{} loaded, err := store.Load(context.Background(), manifestPathFor(cfg)) if err != nil { t.Fatalf("Load(session) error = %v", err) } analyze := loaded.Stages["analyze"] if analyze.Status != manifest.StatusFailed || len(analyze.Outputs) != 0 { t.Fatalf("aggregate analyze state = %#v, want failed without outputs", analyze) } if analyze.AnalyzeArtifacts["player_handout"].Status != manifest.AnalyzeArtifactCurrent || analyze.AnalyzeArtifacts["session_recap"].Status != manifest.AnalyzeArtifactCurrent || analyze.AnalyzeArtifacts["quest_log"].Status != manifest.AnalyzeArtifactFailed { t.Fatalf("partial session projection = %#v", analyze.AnalyzeArtifacts) } if loaded.Stages["publish"].Status != manifest.StatusStale { t.Fatalf("publish status = %q, want stale", loaded.Stages["publish"].Status) } runsDir := artifacts.SessionRunsDirForCampaign(cfg.Pipeline.Workspace.Root, cfg.Session.Campaign, cfg.Session.SessionID) entries, err := os.ReadDir(runsDir) if err != nil || len(entries) != 1 { t.Fatalf("run directory entries = %#v, error = %v", entries, err) } runManifest, err := store.LoadRun(context.Background(), filepath.Join(runsDir, entries[0].Name(), "manifest.json")) if err != nil { t.Fatalf("LoadRun() error = %v", err) } runAnalyze := runManifest.Stages["analyze"] if runAnalyze.Status != manifest.StatusFailed || len(runAnalyze.AnalyzeArtifacts) != 2 || len(runAnalyze.Outputs) != 0 { t.Fatalf("partial invocation projection = %#v", runAnalyze) } } func TestExecuteStagesRejectsInvalidAnalyzeProjectionWithoutReplacingPriorState(t *testing.T) { cfg := testConfig(t) now := time.Date(2026, 5, 3, 9, 0, 0, 0, time.UTC) prior := appAnalyzeRecord("player_handout", manifest.AnalyzeArtifactCurrent, "old-run", now) seed := manifest.New(cfg.Session.SessionID, now) seed.Campaign = cfg.Session.Campaign seed.MarkStageSucceeded("analyze", now, nil) seed.Stages["analyze"].AnalyzeStateVersion = manifest.AnalyzeStateContractVersion seed.Stages["analyze"].AnalyzeArtifacts = map[string]manifest.AnalyzeArtifactRecord{"player_handout": prior} saveBoundedManifest(t, cfg, seed) stageToRun := projectionStage{name: "analyze", run: func(_ *stage.Env, m *manifest.Manifest) (*stage.StageResult, error) { invalid := appAnalyzeRecord("session_recap", manifest.AnalyzeArtifactCurrent, m.RunID, time.Now().UTC()) invalid.Output.Checksum = "invalid" return &stage.StageResult{AnalyzeState: &stage.AnalyzeStateProjection{ Session: map[string]manifest.AnalyzeArtifactRecord{"session_recap": invalid}, Invocation: map[string]manifest.AnalyzeArtifactRecord{"session_recap": invalid}, }}, errors.New("analysis failed") }} _, err := executeStages(context.Background(), cfg, []stage.Stage{stageToRun}, RunOptions{Force: true}) if err == nil || !strings.Contains(err.Error(), "checksum") { t.Fatalf("executeStages() error = %v, want projection validation failure", err) } loaded, loadErr := (&manifest.LocalStore{}).Load(context.Background(), manifestPathFor(cfg)) if loadErr != nil { t.Fatalf("Load() error = %v", loadErr) } if len(loaded.Stages["analyze"].AnalyzeArtifacts) != 1 || !reflect.DeepEqual(loaded.Stages["analyze"].AnalyzeArtifacts["player_handout"], prior) { t.Fatalf("prior state replaced by invalid projection: %#v", loaded.Stages["analyze"].AnalyzeArtifacts) } } func TestExecuteStagesRollsBackAnalyzeAuthorityWhenProjectionSaveFails(t *testing.T) { cfg := testConfig(t) now := time.Date(2026, 5, 3, 9, 0, 0, 0, time.UTC) prior := appAnalyzeRecord("player_handout", manifest.AnalyzeArtifactCurrent, "old-run", now) seed := manifest.New(cfg.Session.SessionID, now) seed.Campaign = cfg.Session.Campaign seed.MarkStageSucceeded("analyze", now, nil) seed.Stages["analyze"].AnalyzeStateVersion = manifest.AnalyzeStateContractVersion seed.Stages["analyze"].AnalyzeArtifacts = map[string]manifest.AnalyzeArtifactRecord{"player_handout": prior} saveBoundedManifest(t, cfg, seed) stageReturned := false store := &analyzeProjectionFailingStore{delegate: &manifest.LocalStore{}, shouldFail: func(m *manifest.Manifest) bool { return stageReturned && m.Stages["analyze"] != nil && m.Stages["analyze"].Status == manifest.StatusSucceeded && m.Stages["analyze"].AnalyzeArtifacts["session_recap"].Status == manifest.AnalyzeArtifactCurrent }} stageToRun := projectionStage{name: "analyze", run: func(_ *stage.Env, m *manifest.Manifest) (*stage.StageResult, error) { stageReturned = true newRecord := appAnalyzeRecord("session_recap", manifest.AnalyzeArtifactCurrent, m.RunID, time.Now().UTC()) return &stage.StageResult{AnalyzeState: &stage.AnalyzeStateProjection{ Session: map[string]manifest.AnalyzeArtifactRecord{ "player_handout": prior, "session_recap": newRecord, }, Invocation: map[string]manifest.AnalyzeArtifactRecord{"session_recap": newRecord}, }}, nil }} _, err := executeStages(context.Background(), cfg, []stage.Stage{stageToRun}, RunOptions{ Force: true, Env: &Env{ManifestStore: store}, }) if err == nil || !strings.Contains(err.Error(), "injected analyze projection save failure") { t.Fatalf("executeStages() error = %v", err) } if !store.failed { t.Fatal("projection persistence failure was not injected") } loaded, loadErr := store.delegate.Load(context.Background(), manifestPathFor(cfg)) if loadErr != nil { t.Fatalf("Load() error = %v", loadErr) } analyze := loaded.Stages["analyze"] if analyze.Status != manifest.StatusFailed || len(analyze.AnalyzeArtifacts) != 1 || !reflect.DeepEqual(analyze.AnalyzeArtifacts["player_handout"], prior) { t.Fatalf("durable analyze state after rollback = %#v", analyze) } } func TestExecuteStagesRejectsAnalyzeProjectionFromOtherStage(t *testing.T) { cfg := testConfig(t) stageToRun := projectionStage{name: "prepare", run: func(_ *stage.Env, m *manifest.Manifest) (*stage.StageResult, error) { record := appAnalyzeRecord("session_recap", manifest.AnalyzeArtifactCurrent, m.RunID, time.Now().UTC()) return &stage.StageResult{AnalyzeState: &stage.AnalyzeStateProjection{ Session: map[string]manifest.AnalyzeArtifactRecord{"session_recap": record}, }}, nil }} _, err := executeStages(context.Background(), cfg, []stage.Stage{stageToRun}, RunOptions{}) if err == nil || !strings.Contains(err.Error(), "returned analyze-owned state projection") { t.Fatalf("executeStages() error = %v", err) } } func TestExecuteStagesRejectsContradictoryAnalyzeResultWithError(t *testing.T) { cfg := testConfig(t) stageToRun := projectionStage{name: "analyze", run: func(_ *stage.Env, m *manifest.Manifest) (*stage.StageResult, error) { record := appAnalyzeRecord("session_recap", manifest.AnalyzeArtifactCurrent, m.RunID, time.Now().UTC()) return &stage.StageResult{ Outputs: []artifacts.Ref{{Kind: "session_recap", RelativePath: "artifacts/session-recap.md"}}, AnalyzeState: &stage.AnalyzeStateProjection{ Session: map[string]manifest.AnalyzeArtifactRecord{"session_recap": record}, Invocation: map[string]manifest.AnalyzeArtifactRecord{"session_recap": record}, }, }, errors.New("analysis failed") }} _, err := executeStages(context.Background(), cfg, []stage.Stage{stageToRun}, RunOptions{}) if err == nil || !strings.Contains(err.Error(), "may contain only analyze-owned state projection") { t.Fatalf("executeStages() error = %v", err) } } func TestExecuteStagesExposesSelectedForceDecisionToStage(t *testing.T) { for _, force := range []bool{false, true} { t.Run(strings.ToLower(strings.TrimSpace(map[bool]string{false: "ordinary", true: "forced"}[force])), func(t *testing.T) { cfg := testConfig(t) captured := !force stageToRun := projectionStage{name: "prepare", run: func(env *stage.Env, _ *manifest.Manifest) (*stage.StageResult, error) { captured = env.Force return &stage.StageResult{}, nil }} if _, err := executeStages(context.Background(), cfg, []stage.Stage{stageToRun}, RunOptions{Force: force}); err != nil { t.Fatalf("executeStages() error = %v", err) } if captured != force { t.Fatalf("stage env force = %v, want %v", captured, force) } }) } } func appAnalyzeRecord(key string, status manifest.AnalyzeArtifactStatus, producerRunID string, at time.Time) manifest.AnalyzeArtifactRecord { record := manifest.AnalyzeArtifactRecord{ Key: key, Status: status, ProducerRunID: producerRunID, UpdatedAt: at, } if status == manifest.AnalyzeArtifactFailed { record.Error = "scriptorium failed" return record } if status != manifest.AnalyzeArtifactCurrent { record.FingerprintVersion = manifest.AnalyzeFingerprintContractVersion record.Fingerprint = strings.Repeat("b", 64) return record } record.FingerprintVersion = manifest.AnalyzeFingerprintContractVersion record.Fingerprint = strings.Repeat("a", 64) record.OutputSize = 42 record.Output = &manifest.ArtifactRecord{ Kind: key, SourceID: artifacts.ConfiguredArtifactSourceID(key), LocalPath: "artifacts/" + strings.ReplaceAll(key, "_", "-") + ".md", ProducerRunID: producerRunID, Checksum: strings.Repeat("c", 64), Contract: &artifactmodel.ContractMetadata{ MediaType: "text/markdown", SchemaID: "narratio." + key, SchemaVersion: "1", }, } return record } func analyzeArtifactOutputKeys(outputs []manifest.ArtifactRecord) []string { keys := make([]string, 0, len(outputs)) for _, output := range outputs { keys = append(keys, strings.TrimPrefix(output.SourceID, "narratio.artifact.")) } return keys } type analyzeProjectionFailingStore struct { delegate *manifest.LocalStore shouldFail func(*manifest.Manifest) bool failed bool } func (s *analyzeProjectionFailingStore) Create(ctx context.Context, sessionID string) (*manifest.Manifest, error) { return s.delegate.Create(ctx, sessionID) } func (s *analyzeProjectionFailingStore) Load(ctx context.Context, path string) (*manifest.Manifest, error) { return s.delegate.Load(ctx, path) } func (s *analyzeProjectionFailingStore) Save(ctx context.Context, path string, m *manifest.Manifest) error { if !s.failed && s.shouldFail != nil && s.shouldFail(m) { s.failed = true return errors.New("injected analyze projection save failure") } return s.delegate.Save(ctx, path, m) }