347 lines
15 KiB
Go
347 lines
15 KiB
Go
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)
|
|
}
|