Files
narratio/internal/app/extract_lifecycle_test.go

734 lines
32 KiB
Go

package app
import (
"context"
"encoding/json"
"errors"
"fmt"
"os"
"path/filepath"
"sort"
"strings"
"testing"
"time"
"gitea.maximumdirect.net/eric/narratio/internal/adapters/notarius"
"gitea.maximumdirect.net/eric/narratio/internal/artifactmodel"
"gitea.maximumdirect.net/eric/narratio/internal/artifactpolicy"
"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"
)
type materializingNotariusRunner struct {
cfg *config.NotariusConfig
requests []notarius.RunRequest
failuresRemaining int
}
type assertExtractionSourcesStage struct {
keys []string
runs *int
}
func (s assertExtractionSourcesStage) Name() string { return "analyze" }
func (s assertExtractionSourcesStage) Run(_ context.Context, env *stage.Env, m *manifest.Manifest) (*stage.StageResult, error) {
definitions := artifacts.ExtractionDefinitionsFromConfig(env.Config.Pipeline.Notarius)
catalog, err := artifacts.BootstrapRuntimeCatalog(nil, env.EffectiveArtifacts, definitions)
if err != nil {
return nil, err
}
paths, err := env.ArtifactStore.EnsureLayoutFor(env.Config.Session.Campaign, env.Config.Session.SessionID)
if err != nil {
return nil, err
}
catalog.HydrateExtractionArtifacts(paths, m, definitions)
for _, key := range s.keys {
sourceID := artifacts.ExtractionArtifactSourceID(key)
entry, ok := catalog.Lookup(sourceID)
if !ok || !entry.Available || entry.SourceID != sourceID || entry.Path == "" {
return nil, fmt.Errorf("extraction source %q unavailable: %#v, present=%v", sourceID, entry, ok)
}
}
if s.runs != nil {
*s.runs = *s.runs + 1
}
return &stage.StageResult{}, nil
}
func (r *materializingNotariusRunner) Run(_ context.Context, req notarius.RunRequest) (notarius.RunResult, error) {
r.requests = append(r.requests, req)
if r.failuresRemaining > 0 {
r.failuresRemaining--
return notarius.RunResult{}, errors.New("notarius execution failed")
}
externalRunID := fmt.Sprintf("notarius-run-%d", len(r.requests))
bundle := filepath.Join(req.OutputRoot, externalRunID)
lanesDir := filepath.Join(bundle, "lanes")
if err := os.MkdirAll(lanesDir, 0o755); err != nil {
return notarius.RunResult{}, err
}
for path, content := range map[string]string{
filepath.Join(bundle, "index.json"): `{"manifest_file":"manifest.json"}`,
filepath.Join(bundle, "manifest.json"): `{}`,
filepath.Join(bundle, "rejected.json"): `{"rejected":[]}`,
filepath.Join(bundle, "warnings.json"): `{"schema_version":"notarius.warnings.v2","group_count":0,"occurrence_count":0,"groups":[]}`,
filepath.Join(bundle, "diagnostics.json"): `{"schema_version":"notarius.diagnostics.v1","group_count":0,"occurrence_count":0,"truncated":false,"unrepresented_occurrence_count":0,"groups":[]}`,
} {
if err := os.WriteFile(path, []byte(content), 0o644); err != nil {
return notarius.RunResult{}, err
}
}
keys := make([]string, 0, len(r.cfg.Outputs))
for key := range r.cfg.Outputs {
keys = append(keys, key)
}
sort.Strings(keys)
lanes := make([]notarius.LaneDescriptor, 0, len(keys))
for _, key := range keys {
output := r.cfg.Outputs[key]
filename := key + ".json"
path := filepath.Join(lanesDir, filename)
if err := os.WriteFile(path, []byte(`{"records":[]}`), 0o644); err != nil {
return notarius.RunResult{}, err
}
lanes = append(lanes, notarius.LaneDescriptor{
LaneID: output.LaneID, File: filepath.ToSlash(filepath.Join("lanes", filename)), Path: path,
MediaType: output.MediaType, SchemaID: output.SchemaID,
SchemaVersion: output.SchemaVersion, ModuleKey: output.ModuleKey,
})
}
return notarius.RunResult{
Receipt: notarius.Receipt{
SchemaVersion: notarius.ReceiptSchemaVersion, RunID: externalRunID,
PipelineID: req.PipelineID, OutputDirectory: bundle, IndexFile: "index.json",
NormalizedOutputCount: len(lanes), ValidationStatus: "approved",
},
BundleRoot: bundle,
Index: notarius.Index{
Path: filepath.Join(bundle, "index.json"), RejectedPath: filepath.Join(bundle, "rejected.json"),
WarningsPath: filepath.Join(bundle, "warnings.json"),
DiagnosticsPath: filepath.Join(bundle, "diagnostics.json"),
Lanes: lanes,
},
}, nil
}
func TestExtractLifecycleDisabledThenEnabled(t *testing.T) {
cfg, env, runner := extractionLifecycleFixture(t, false)
plan, err := BuildSingleStagePlan("extract")
if err != nil {
t.Fatalf("BuildSingleStagePlan(extract) error = %v", err)
}
first, err := executeStages(context.Background(), cfg, plan, RunOptions{Env: env})
if err != nil {
t.Fatalf("disabled executeStages() error = %v", err)
}
if len(first.Executed) != 1 || len(first.Skipped) != 1 || first.Skipped[0] != "extract" || len(runner.requests) != 0 {
t.Fatalf("disabled summary = %#v requests=%d", first, len(runner.requests))
}
cfg.Pipeline.Notarius.Enabled = true
second, err := executeStages(context.Background(), cfg, plan, RunOptions{Env: env})
if err != nil {
t.Fatalf("enabled executeStages() error = %v", err)
}
if len(second.Executed) != 1 || len(second.Skipped) != 0 || len(runner.requests) != 1 {
t.Fatalf("enabled summary = %#v requests=%d", second, len(runner.requests))
}
loaded, err := (&manifest.LocalStore{}).Load(context.Background(), second.ManifestPath)
if err != nil {
t.Fatalf("Load() error = %v", err)
}
if loaded.Stages["extract"] == nil || loaded.Stages["extract"].Status != manifest.StatusSucceeded || len(loaded.Stages["extract"].Outputs) != 2 {
t.Fatalf("extract record = %#v, want succeeded manifest-ready outputs", loaded.Stages["extract"])
}
}
func TestExtractLifecyclePrepareBindsVerifiedReferenceSnapshots(t *testing.T) {
cfg, env, runner := extractionLifecycleFixture(t, true)
originalPaths := configureLifecycleReferences(t, cfg)
summary, err := executeStages(context.Background(), cfg, prepareExtractLifecyclePlan(t), RunOptions{Env: env})
if err != nil {
t.Fatalf("executeStages() error = %v", err)
}
if len(summary.Executed) != 2 || len(runner.requests) != 1 {
t.Fatalf("summary = %#v requests=%d", summary, len(runner.requests))
}
paths, err := artifacts.NewLocalStore(cfg.Pipeline.Workspace.Root).EnsureLayoutFor(cfg.Session.Campaign, cfg.Session.SessionID)
if err != nil {
t.Fatalf("EnsureLayoutFor() error = %v", err)
}
want := []struct {
selector string
sourceID string
filename string
}{
{selector: "glossary", sourceID: artifactpolicy.SourceInputGlossary, filename: "glossary.yml"},
{selector: "party", sourceID: artifactpolicy.SourceInputParty, filename: "party.yml"},
{selector: "players", sourceID: artifactpolicy.SourceInputPlayers, filename: "players.yml"},
{selector: "spells", sourceID: artifactpolicy.SourceInputSpellCatalog, filename: "spell_catalog.json"},
}
request := runner.requests[0]
if len(request.References) != len(want) {
t.Fatalf("references = %#v", request.References)
}
loaded := loadLifecycleManifest(t, cfg)
for index, expected := range want {
binding := request.References[index]
canonical := filepath.Join(paths.InputsDir, expected.filename)
snapshot := filepath.Join(
artifacts.SessionRunNotariusReferencesDirForCampaign(
cfg.Pipeline.Workspace.Root, cfg.Session.Campaign, cfg.Session.SessionID, loaded.RunID,
),
expected.filename,
)
if binding.Selector != expected.selector || binding.Path != snapshot || binding.Path == canonical || binding.Path == originalPaths[expected.sourceID] {
t.Fatalf("reference[%d] = %#v, want selector %q snapshot %q and not prepared/source paths", index, binding, expected.selector, snapshot)
}
}
extract := loaded.Stages["extract"]
if extract == nil || extract.Status != manifest.StatusSucceeded || extract.Metadata["reference_count"] != float64(len(want)) {
t.Fatalf("extract record = %#v", extract)
}
references, ok := extract.Metadata["references"].([]any)
if !ok || len(references) != len(want) || len(references) > config.MaxNotariusReferenceBindings {
t.Fatalf("reference metadata = %#v", extract.Metadata["references"])
}
for index, raw := range references {
entry, ok := raw.(map[string]any)
if !ok || len(entry) != 5 || entry["selector"] != want[index].selector || entry["source_id"] != want[index].sourceID {
t.Fatalf("reference metadata[%d] = %#v", index, raw)
}
}
}
func TestExtractLifecyclePreparedReferenceChangeRerunsExtractionAndInvalidatesDownstream(t *testing.T) {
cfg, env, runner := extractionLifecycleFixture(t, true)
originalPaths := configureLifecycleReferences(t, cfg)
if _, err := executeStages(context.Background(), cfg, prepareExtractLifecyclePlan(t), RunOptions{Env: env}); err != nil {
t.Fatalf("initial executeStages() error = %v", err)
}
before := loadLifecycleManifest(t, cfg)
beforeChecksum := lifecycleInputChecksum(t, before, "party")
for _, name := range []string{"render", "analyze", "publish"} {
before.MarkStageSucceeded(name, time.Now().UTC(), nil)
}
if err := (&manifest.LocalStore{}).Save(context.Background(), manifestPathFor(cfg), before); err != nil {
t.Fatalf("Save(downstream success) error = %v", err)
}
if err := os.WriteFile(originalPaths[artifactpolicy.SourceInputParty], []byte("changed party bytes\n"), 0o644); err != nil {
t.Fatalf("WriteFile(party source) error = %v", err)
}
prepared := loadLifecycleManifest(t, cfg)
prepare, err := stage.Select("prepare")
if err != nil {
t.Fatalf("stage.Select(prepare) error = %v", err)
}
if _, err := prepare.Run(context.Background(), env, prepared); err != nil {
t.Fatalf("prepare.Run() error = %v", err)
}
if lifecycleInputChecksum(t, prepared, "party") == beforeChecksum {
t.Fatal("prepared party checksum did not change")
}
if err := (&manifest.LocalStore{}).Save(context.Background(), manifestPathFor(cfg), prepared); err != nil {
t.Fatalf("Save(reprepared manifest) error = %v", err)
}
plan, err := BuildSingleStagePlan("extract")
if err != nil {
t.Fatalf("BuildSingleStagePlan(extract) error = %v", err)
}
run, err := executeStages(context.Background(), cfg, plan, RunOptions{Env: env})
if err != nil {
t.Fatalf("rerun executeStages() error = %v", err)
}
if len(run.Executed) != 1 || len(run.Skipped) != 0 || len(runner.requests) != 2 {
t.Fatalf("rerun summary = %#v requests=%d", run, len(runner.requests))
}
after := loadLifecycleManifest(t, cfg)
if after.Stages["extract"].Status != manifest.StatusSucceeded {
t.Fatalf("extract status = %#v", after.Stages["extract"])
}
for _, name := range []string{"render", "analyze", "publish"} {
if after.Stages[name] == nil || after.Stages[name].Status != manifest.StatusStale {
t.Fatalf("%s status = %#v, want stale", name, after.Stages[name])
}
}
}
func TestExtractLifecycleSessionOverrideBytesReachCanonicalReference(t *testing.T) {
cfg, env, runner := extractionLifecycleFixture(t, true)
configureLifecycleReferences(t, cfg)
overridePath := filepath.Join(filepath.Dir(cfg.SessionPath), "session-party.yml")
if err := os.WriteFile(overridePath, []byte("session override party\n"), 0o644); err != nil {
t.Fatalf("WriteFile(session override) error = %v", err)
}
cfg.StableInputs.PartyFile = config.ResolvedInputFile{
Path: "./session-party.yml", ConfigPath: cfg.SessionPath, Source: "session_config",
}
cfg.Session.Inputs.PartyFile = "./session-party.yml"
if _, err := executeStages(context.Background(), cfg, prepareExtractLifecyclePlan(t), RunOptions{Env: env}); err != nil {
t.Fatalf("executeStages() error = %v", err)
}
if len(runner.requests) != 1 {
t.Fatalf("requests = %d", len(runner.requests))
}
var partyPath string
for _, binding := range runner.requests[0].References {
if binding.Selector == "party" {
partyPath = binding.Path
}
}
contents, err := os.ReadFile(partyPath)
if err != nil {
t.Fatalf("ReadFile(prepared party) error = %v", err)
}
if string(contents) != "session override party\n" || partyPath == overridePath {
t.Fatalf("prepared party path=%q contents=%q override=%q", partyPath, contents, overridePath)
}
}
func TestExtractLifecycleEmptyReferencesPreserveAllDndExtractionSources(t *testing.T) {
cfg, env, runner := extractionLifecycleFixture(t, true)
cfg.Pipeline.Notarius.Outputs = lifecycleDndOutputs()
keys := make([]string, 0, len(cfg.Pipeline.Notarius.Outputs))
for key := range cfg.Pipeline.Notarius.Outputs {
keys = append(keys, key)
}
sort.Strings(keys)
analyzeRuns := 0
extractPlan, err := BuildSingleStagePlan("extract")
if err != nil {
t.Fatalf("BuildSingleStagePlan(extract) error = %v", err)
}
plan := append(extractPlan, assertExtractionSourcesStage{keys: keys, runs: &analyzeRuns})
summary, err := executeStages(context.Background(), cfg, plan, RunOptions{Env: env})
if err != nil {
t.Fatalf("executeStages() error = %v", err)
}
if len(summary.Executed) != 2 || len(runner.requests) != 1 || len(runner.requests[0].References) != 0 || analyzeRuns != 1 {
t.Fatalf("summary=%#v requests=%#v analyze=%d", summary, runner.requests, analyzeRuns)
}
loaded := loadLifecycleManifest(t, cfg)
if got := len(loaded.Stages["extract"].Outputs); got != len(keys)+1 {
t.Fatalf("extract outputs = %d, want %d lanes plus index", got, len(keys))
}
}
func TestExtractLifecycleChangedOutcomeRerunsSucceededDownstream(t *testing.T) {
cfg, env, runner := extractionLifecycleFixture(t, false)
analyzeRuns := 0
plan := extractionLifecyclePlan(t, &analyzeRuns)
if _, err := executeStages(context.Background(), cfg, plan, RunOptions{Env: env}); err != nil {
t.Fatalf("disabled executeStages() error = %v", err)
}
if analyzeRuns != 1 || len(runner.requests) != 0 {
t.Fatalf("disabled run analyze=%d Notarius=%d, want 1 and 0", analyzeRuns, len(runner.requests))
}
cfg.Pipeline.Notarius.Enabled = true
summary, err := executeStages(context.Background(), cfg, plan, RunOptions{Env: env})
if err != nil {
t.Fatalf("enabled executeStages() error = %v", err)
}
if analyzeRuns != 2 || len(runner.requests) != 1 {
t.Fatalf("enabled run analyze=%d Notarius=%d, want 2 and 1", analyzeRuns, len(runner.requests))
}
if len(summary.Executed) != 2 || len(summary.Skipped) != 0 {
t.Fatalf("enabled summary = %#v, want extract and analyze executed", summary)
}
}
func TestExtractLifecycleFailureInvalidatesAndOrdinaryRetryRerunsDownstream(t *testing.T) {
cfg, env, runner := extractionLifecycleFixture(t, true)
runner.failuresRemaining = 1
markLifecycleStageSucceeded(t, cfg, "analyze")
analyzeRuns := 0
plan := extractionLifecyclePlan(t, &analyzeRuns)
if _, err := executeStages(context.Background(), cfg, plan, RunOptions{Env: env}); err == nil || !strings.Contains(err.Error(), "notarius execution failed") {
t.Fatalf("failed executeStages() error = %v", err)
}
failed := loadLifecycleManifest(t, cfg)
if failed.Stages["extract"].Status != manifest.StatusFailed || failed.Stages["analyze"].Status != manifest.StatusStale {
t.Fatalf("failed lifecycle extract=%#v analyze=%#v", failed.Stages["extract"], failed.Stages["analyze"])
}
if failed.Stages["analyze"].Error == nil || failed.Stages["analyze"].Error.Message != staleReasonFailure {
t.Fatalf("analyze stale reason = %#v, want %q", failed.Stages["analyze"].Error, staleReasonFailure)
}
summary, err := executeStages(context.Background(), cfg, plan, RunOptions{Env: env})
if err != nil {
t.Fatalf("retry executeStages() error = %v", err)
}
if analyzeRuns != 1 || len(runner.requests) != 2 || len(summary.Executed) != 2 {
t.Fatalf("retry analyze=%d Notarius=%d summary=%#v", analyzeRuns, len(runner.requests), summary)
}
}
func TestExtractLifecycleForcedSelfSkipInvalidatesDownstream(t *testing.T) {
cfg, env, _ := extractionLifecycleFixture(t, false)
markLifecycleStageSucceeded(t, cfg, "analyze")
plan, _ := BuildSingleStagePlan("extract")
if _, err := executeStages(context.Background(), cfg, plan, RunOptions{Env: env, Force: true}); err != nil {
t.Fatalf("executeStages() error = %v", err)
}
loaded := loadLifecycleManifest(t, cfg)
if loaded.Stages["extract"].Status != manifest.StatusSkipped || loaded.Stages["analyze"].Status != manifest.StatusStale {
t.Fatalf("forced self-skip extract=%#v analyze=%#v", loaded.Stages["extract"], loaded.Stages["analyze"])
}
if loaded.Stages["analyze"].Error == nil || loaded.Stages["analyze"].Error.Message != staleReasonForcedReplacement {
t.Fatalf("analyze stale reason = %#v, want %q", loaded.Stages["analyze"].Error, staleReasonForcedReplacement)
}
}
func TestExtractLifecycleForcedFailureInvalidatesDownstream(t *testing.T) {
cfg, env, runner := extractionLifecycleFixture(t, true)
plan, _ := BuildSingleStagePlan("extract")
first, err := executeStages(context.Background(), cfg, plan, RunOptions{Env: env})
if err != nil {
t.Fatalf("initial executeStages() error = %v", err)
}
succeeded := loadLifecycleManifest(t, cfg).Stages["extract"]
if succeeded == nil || succeeded.Status != manifest.StatusSucceeded || len(succeeded.Outputs) == 0 || len(succeeded.Logs) == 0 || len(succeeded.Metadata) == 0 {
t.Fatalf("initial extraction result = %#v, want succeeded result details", succeeded)
}
markLifecycleStageSucceeded(t, cfg, "analyze")
runner.failuresRemaining = 1
if _, err := executeStages(context.Background(), cfg, plan, RunOptions{Env: env, Force: true}); err == nil {
t.Fatal("executeStages() error = nil, want forced extraction failure")
}
loaded := loadLifecycleManifest(t, cfg)
if loaded.Stages["extract"].Status != manifest.StatusFailed || loaded.Stages["analyze"].Status != manifest.StatusStale {
t.Fatalf("forced failure extract=%#v analyze=%#v", loaded.Stages["extract"], loaded.Stages["analyze"])
}
if loaded.Stages["analyze"].Error == nil || loaded.Stages["analyze"].Error.Message != staleReasonForcedReplacement {
t.Fatalf("analyze stale reason = %#v, want %q", loaded.Stages["analyze"].Error, staleReasonForcedReplacement)
}
failed := loaded.Stages["extract"]
if len(failed.Outputs) != 0 || len(failed.Logs) != 0 || len(failed.GeneratedConfigs) != 0 || len(failed.Metadata) != 0 {
t.Fatalf("failed replacement inherited extraction result details: %#v", failed)
}
historical, err := (&manifest.LocalStore{}).LoadRun(context.Background(), first.RunManifestPath)
if err != nil {
t.Fatalf("LoadRun(initial) error = %v", err)
}
historicalExtract := historical.Stages["extract"]
if historicalExtract == nil || historicalExtract.Status != manifest.StatusSucceeded || len(historicalExtract.Outputs) == 0 || len(historicalExtract.Logs) == 0 || len(historicalExtract.Metadata) == 0 {
t.Fatalf("historical extraction result = %#v, want preserved succeeded details", historicalExtract)
}
if _, err := os.Stat(succeeded.Outputs[0].LocalPath); err != nil {
t.Fatalf("durable extraction output was not preserved: %v", err)
}
}
func TestExtractLifecycleRepeatedSelfSkipPreservesSucceededDownstream(t *testing.T) {
cfg, env, runner := extractionLifecycleFixture(t, false)
analyzeRuns := 0
plan := extractionLifecyclePlan(t, &analyzeRuns)
if _, err := executeStages(context.Background(), cfg, plan, RunOptions{Env: env}); err != nil {
t.Fatalf("first executeStages() error = %v", err)
}
second, err := executeStages(context.Background(), cfg, plan, RunOptions{Env: env})
if err != nil {
t.Fatalf("second executeStages() error = %v", err)
}
if analyzeRuns != 1 || len(runner.requests) != 0 {
t.Fatalf("repeated disabled run analyze=%d Notarius=%d, want 1 and 0", analyzeRuns, len(runner.requests))
}
if len(second.Executed) != 1 || len(second.Skipped) != 2 {
t.Fatalf("second summary = %#v, want executed self-skip and skipped analyze", second)
}
loaded := loadLifecycleManifest(t, cfg)
if loaded.Stages["analyze"].Status != manifest.StatusSucceeded {
t.Fatalf("analyze = %#v, want succeeded", loaded.Stages["analyze"])
}
}
func TestExtractLifecycleSkipsCurrentResumableResult(t *testing.T) {
cfg, env, runner := extractionLifecycleFixture(t, true)
analyzeRuns := 0
plan := extractionLifecyclePlan(t, &analyzeRuns)
if _, err := executeStages(context.Background(), cfg, plan, RunOptions{Env: env}); err != nil {
t.Fatalf("first executeStages() error = %v", err)
}
resumed, err := executeStages(context.Background(), cfg, plan, RunOptions{Env: env})
if err != nil {
t.Fatalf("resume executeStages() error = %v", err)
}
if len(resumed.Executed) != 0 || len(resumed.Skipped) != 2 || len(runner.requests) != 1 || analyzeRuns != 1 {
t.Fatalf("resume summary = %#v requests=%d analyze=%d", resumed, len(runner.requests), analyzeRuns)
}
}
func TestExtractLifecycleRerunsAfterDirectTranscriptChange(t *testing.T) {
cfg, env, runner := extractionLifecycleFixture(t, true)
analyzeRuns := 0
plan := extractionLifecyclePlan(t, &analyzeRuns)
if _, err := executeStages(context.Background(), cfg, plan, RunOptions{Env: env}); err != nil {
t.Fatalf("first executeStages() error = %v", err)
}
persisted := loadLifecycleManifest(t, cfg)
trimmed := persisted.Stages["trim"].Outputs[0].LocalPath
if err := os.WriteFile(trimmed, []byte(`{"segments":[{"id":"changed"}]}`), 0o644); err != nil {
t.Fatalf("WriteFile(trimmed transcript) error = %v", err)
}
if err := (&manifest.LocalStore{}).Save(context.Background(), manifestPathFor(cfg), persisted); err != nil {
t.Fatalf("Save(mutated) error = %v", err)
}
rerun, err := executeStages(context.Background(), cfg, plan, RunOptions{Env: env})
if err != nil {
t.Fatalf("rerun executeStages() error = %v", err)
}
if len(rerun.Executed) != 2 || len(rerun.Skipped) != 0 || len(runner.requests) != 2 || analyzeRuns != 2 {
t.Fatalf("rerun summary = %#v requests=%d analyze=%d", rerun, len(runner.requests), analyzeRuns)
}
}
func TestExtractLifecycleResumesAndRerunsObsoleteResults(t *testing.T) {
for _, test := range []struct {
name string
mutate func(*testing.T, *config.Config, *manifest.Manifest)
}{
{name: "configuration changed", mutate: func(_ *testing.T, cfg *config.Config, _ *manifest.Manifest) {
output := cfg.Pipeline.Notarius.Outputs["npc_registry"]
output.SchemaVersion = "v2"
cfg.Pipeline.Notarius.Outputs["npc_registry"] = output
}},
{name: "payload missing", mutate: func(t *testing.T, _ *config.Config, m *manifest.Manifest) {
if err := os.Remove(m.Stages["extract"].Outputs[0].LocalPath); err != nil {
t.Fatalf("Remove() error = %v", err)
}
}},
{name: "payload tampered", mutate: func(t *testing.T, _ *config.Config, m *manifest.Manifest) {
if err := os.WriteFile(m.Stages["extract"].Outputs[0].LocalPath, []byte(`{"npcs":["tampered"]}`), 0o644); err != nil {
t.Fatalf("WriteFile() error = %v", err)
}
}},
{name: "record incompatible", mutate: func(_ *testing.T, _ *config.Config, m *manifest.Manifest) {
m.Stages["extract"].Outputs[0].Contract.SchemaID = "incompatible"
}},
} {
t.Run(test.name, func(t *testing.T) {
cfg, env, runner := extractionLifecycleFixture(t, true)
plan, _ := BuildSingleStagePlan("extract")
first, err := executeStages(context.Background(), cfg, plan, RunOptions{Env: env})
if err != nil {
t.Fatalf("first executeStages() error = %v", err)
}
persisted, err := (&manifest.LocalStore{}).Load(context.Background(), first.ManifestPath)
if err != nil {
t.Fatalf("Load() error = %v", err)
}
test.mutate(t, cfg, persisted)
if err := (&manifest.LocalStore{}).Save(context.Background(), first.ManifestPath, persisted); err != nil {
t.Fatalf("Save(mutated) error = %v", err)
}
rerun, err := executeStages(context.Background(), cfg, plan, RunOptions{Env: env})
if err != nil {
t.Fatalf("rerun executeStages() error = %v", err)
}
if len(rerun.Executed) != 1 || len(rerun.Skipped) != 0 || len(runner.requests) != 2 {
t.Fatalf("rerun summary = %#v requests=%d", rerun, len(runner.requests))
}
})
}
}
func TestExtractLifecycleUnsafeResumeErrorPreservesSuccess(t *testing.T) {
cfg, env, runner := extractionLifecycleFixture(t, true)
analyzeRuns := 0
plan := extractionLifecyclePlan(t, &analyzeRuns)
first, err := executeStages(context.Background(), cfg, plan, RunOptions{Env: env})
if err != nil {
t.Fatalf("first executeStages() error = %v", err)
}
store := &manifest.LocalStore{}
persisted, err := store.Load(context.Background(), first.ManifestPath)
if err != nil {
t.Fatalf("Load() error = %v", err)
}
persisted.Stages["extract"].Outputs[0].LocalPath = filepath.Join(cfg.Pipeline.Workspace.Root, "outside.json")
if err := store.Save(context.Background(), first.ManifestPath, persisted); err != nil {
t.Fatalf("Save() error = %v", err)
}
before, _ := json.Marshal(map[string]*manifest.StageRecord{
"extract": persisted.Stages["extract"],
"analyze": persisted.Stages["analyze"],
})
if _, err := executeStages(context.Background(), cfg, plan, RunOptions{Env: env}); err == nil || !strings.Contains(err.Error(), "unsafe") {
t.Fatalf("executeStages() error = %v, want unsafe resume failure", err)
}
afterManifest, err := store.Load(context.Background(), first.ManifestPath)
if err != nil {
t.Fatalf("Load(after) error = %v", err)
}
after, _ := json.Marshal(map[string]*manifest.StageRecord{
"extract": afterManifest.Stages["extract"],
"analyze": afterManifest.Stages["analyze"],
})
if string(before) != string(after) || len(runner.requests) != 1 || analyzeRuns != 1 {
t.Fatalf("successful records changed: before=%s after=%s requests=%d analyze=%d", before, after, len(runner.requests), analyzeRuns)
}
}
func extractionLifecyclePlan(t *testing.T, analyzeRuns *int) []stage.Stage {
t.Helper()
plan, err := BuildSingleStagePlan("extract")
if err != nil {
t.Fatalf("BuildSingleStagePlan(extract) error = %v", err)
}
return append(plan, countingStage{name: "analyze", runs: analyzeRuns})
}
func loadLifecycleManifest(t *testing.T, cfg *config.Config) *manifest.Manifest {
t.Helper()
loaded, err := (&manifest.LocalStore{}).Load(context.Background(), manifestPathFor(cfg))
if err != nil {
t.Fatalf("Load() error = %v", err)
}
return loaded
}
func markLifecycleStageSucceeded(t *testing.T, cfg *config.Config, name string) {
t.Helper()
loaded := loadLifecycleManifest(t, cfg)
loaded.MarkStageSucceeded(name, time.Now().UTC(), nil)
if err := (&manifest.LocalStore{}).Save(context.Background(), manifestPathFor(cfg), loaded); err != nil {
t.Fatalf("Save() error = %v", err)
}
}
func prepareExtractLifecyclePlan(t *testing.T) []stage.Stage {
t.Helper()
prepare, err := BuildSingleStagePlan("prepare")
if err != nil {
t.Fatalf("BuildSingleStagePlan(prepare) error = %v", err)
}
extract, err := BuildSingleStagePlan("extract")
if err != nil {
t.Fatalf("BuildSingleStagePlan(extract) error = %v", err)
}
return append(prepare, extract...)
}
func configureLifecycleReferences(t *testing.T, cfg *config.Config) map[string]string {
t.Helper()
cfg.Pipeline.Notarius.References = map[string]string{
"party": artifactpolicy.SourceInputParty,
"players": artifactpolicy.SourceInputPlayers,
"glossary": artifactpolicy.SourceInputGlossary,
"spells": artifactpolicy.SourceInputSpellCatalog,
}
spellPath := filepath.Join(filepath.Dir(cfg.CampaignPath), "spells.json")
if err := os.WriteFile(spellPath, []byte(`{"spells":[]}`+"\n"), 0o644); err != nil {
t.Fatalf("WriteFile(spell catalog) error = %v", err)
}
cfg.StableInputs.SpellCatalogFile = config.ResolvedInputFile{
Path: "./spells.json", ConfigPath: cfg.CampaignPath, Source: "campaign_config",
}
cfg.Session.Inputs.SpellCatalogFile = "./spells.json"
return map[string]string{
artifactpolicy.SourceInputParty: filepath.Join(filepath.Dir(cfg.CampaignPath), "party.yml"),
artifactpolicy.SourceInputPlayers: filepath.Join(filepath.Dir(cfg.CampaignPath), "players.yml"),
artifactpolicy.SourceInputGlossary: filepath.Join(filepath.Dir(cfg.CampaignPath), "glossary.yml"),
artifactpolicy.SourceInputSpellCatalog: spellPath,
}
}
func lifecycleInputChecksum(t *testing.T, m *manifest.Manifest, kind string) string {
t.Helper()
for _, input := range m.Inputs {
if input.Kind == kind {
if strings.TrimSpace(input.Checksum) == "" {
t.Fatalf("input %q has no checksum: %#v", kind, input)
}
return input.Checksum
}
}
t.Fatalf("manifest input %q not found: %#v", kind, m.Inputs)
return ""
}
func lifecycleDndOutputs() map[string]config.NotariusOutputConfig {
return map[string]config.NotariusOutputConfig{
"item_registry": {LaneID: "item-registry", MediaType: "application/json", SchemaID: "notarius.dnd.item_registry", SchemaVersion: "v1", ModuleKey: "dnd/item-registry"},
"npc_registry": {LaneID: "npc-registry", MediaType: "application/json", SchemaID: "notarius.dnd.npc_registry", SchemaVersion: "v1", ModuleKey: "dnd/npc-registry"},
"location_registry": {LaneID: "location-registry", MediaType: "application/json", SchemaID: "notarius.dnd.location_registry", SchemaVersion: "v1", ModuleKey: "dnd/location-registry"},
"scene_descriptions": {LaneID: "scene-descriptions", MediaType: "application/json", SchemaID: "notarius.dnd.scene_descriptions", SchemaVersion: "v1", ModuleKey: "dnd/scene-descriptions"},
"item_occurrences": {LaneID: "item-occurrences", MediaType: "application/json", SchemaID: "notarius.dnd.item_occurrences", SchemaVersion: "v1", ModuleKey: "dnd/item-occurrences"},
"spells": {LaneID: "spells", MediaType: "application/json", SchemaID: "notarius.dnd.spells", SchemaVersion: "v1", ModuleKey: "dnd/spells"},
"combat_turns": {LaneID: "combat-turns", MediaType: "application/json", SchemaID: "notarius.dnd.combat_turns", SchemaVersion: "v1", ModuleKey: "dnd/combat-turns"},
"npc_occurrences": {LaneID: "npc-occurrences", MediaType: "application/json", SchemaID: "notarius.dnd.npc_occurrences", SchemaVersion: "v1", ModuleKey: "dnd/npc-occurrences"},
"location_occurrences": {LaneID: "location-occurrences", MediaType: "application/json", SchemaID: "notarius.dnd.location_occurrences", SchemaVersion: "v1", ModuleKey: "dnd/location-occurrences"},
"enemy_events": {LaneID: "enemy-events", MediaType: "application/json", SchemaID: "notarius.dnd.enemy_events", SchemaVersion: "v1", ModuleKey: "dnd/enemy-events"},
}
}
func extractionLifecycleFixture(t *testing.T, enabled bool) (*config.Config, *stage.Env, *materializingNotariusRunner) {
t.Helper()
cfg := testConfig(t)
root := cfg.Pipeline.Workspace.Root
binary := filepath.Join(root, "notarius")
configPath := filepath.Join(root, "notarius.yml")
workingDirectory := filepath.Join(root, "notarius-work")
if err := os.WriteFile(binary, []byte("#!/bin/sh\nexit 0\n"), 0o755); err != nil {
t.Fatalf("WriteFile(binary) error = %v", err)
}
if err := os.WriteFile(configPath, []byte("pipelines: {}\n"), 0o644); err != nil {
t.Fatalf("WriteFile(config) error = %v", err)
}
if err := os.Mkdir(workingDirectory, 0o755); err != nil {
t.Fatalf("Mkdir(working directory) error = %v", err)
}
cfg.Pipeline.Notarius = &config.NotariusConfig{
Enabled: enabled, Binary: binary, ConfigPath: configPath, PipelineID: "dnd-session",
Timeout: "45m", WorkingDirectory: workingDirectory,
Outputs: map[string]config.NotariusOutputConfig{
"npc_registry": {
LaneID: "npc-registry", MediaType: "application/json", SchemaID: "notarius.dnd.npc_registry",
SchemaVersion: "v1", ModuleKey: "dnd/npc-registry",
},
},
}
paths, err := artifacts.NewLocalStore(root).EnsureLayoutFor(cfg.Session.Campaign, cfg.Session.SessionID)
if err != nil {
t.Fatalf("EnsureLayoutFor() error = %v", err)
}
inputPath := filepath.Join(paths.ArtifactsDir, "trimmed.from-manifest.json")
if err := os.WriteFile(inputPath, []byte(`{"segments":[]}`), 0o644); err != nil {
t.Fatalf("WriteFile(input) error = %v", err)
}
m := manifest.New(cfg.Session.SessionID, time.Now().UTC())
m.Campaign = cfg.Session.Campaign
m.MarkStageSucceeded("trim", time.Now().UTC(), []manifest.ArtifactRecord{{
Kind: artifactmodel.TranscriptOutputKindFinalTrimmed, SourceID: artifactmodel.SourceTranscriptFinalTrimmed,
LocalPath: inputPath,
}})
if err := (&manifest.LocalStore{}).Save(context.Background(), paths.ManifestPath, m); err != nil {
t.Fatalf("Save(seed) error = %v", err)
}
runner := &materializingNotariusRunner{cfg: cfg.Pipeline.Notarius}
return cfg, &stage.Env{Notarius: runner}, runner
}