diff --git a/docs/roadmap/implementation.md b/docs/roadmap/implementation.md index 032a902..fe301e1 100644 --- a/docs/roadmap/implementation.md +++ b/docs/roadmap/implementation.md @@ -23,7 +23,7 @@ implementation sequence. | Stage 5 | Complete | | Stage 6 | Complete | | Stage 7 | Complete | -| Stage 8 | Not started | +| Stage 8 | Complete | | Stage 9 | Not started | After completing and validating a stage, update only that stage's row to diff --git a/internal/app/operator_artifact_rendering.go b/internal/app/operator_artifact_rendering.go index 0b479d3..687d83b 100644 --- a/internal/app/operator_artifact_rendering.go +++ b/internal/app/operator_artifact_rendering.go @@ -10,9 +10,10 @@ import ( "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" ) -func buildHelperArtifactCatalog(cfg *config.Config) (*artifacts.ArtifactCatalog, error) { +func buildHelperArtifactCatalog(cfg *config.Config, m *manifest.Manifest) (*artifacts.ArtifactCatalog, error) { catalog := artifacts.NewArtifactCatalog() if err := catalog.RegisterBuiltIns(); err != nil { return nil, err @@ -26,6 +27,14 @@ func buildHelperArtifactCatalog(cfg *config.Config) (*artifacts.ArtifactCatalog, if err := catalog.RegisterConfiguredArtifacts(configured, nil); err != nil { return nil, err } + extractionDefinitions := artifacts.ExtractionDefinitionsFromConfig(cfg.Pipeline.Notarius) + if err := catalog.RegisterExtractionArtifacts(extractionDefinitions); err != nil { + return nil, err + } + paths := artifacts.NewLocalStore(cfg.Pipeline.Workspace.Root).SessionPathsFor(cfg.Session.Campaign, cfg.Session.SessionID) + if cfg.Pipeline.Notarius != nil && cfg.Pipeline.Notarius.Enabled { + catalog.HydrateExtractionArtifacts(paths, m, extractionDefinitions) + } return catalog, nil } @@ -40,6 +49,14 @@ func writeArtifactList(out io.Writer, cfg *config.Config, catalog *artifacts.Art for _, entry := range catalog.ListConfigured() { writeArtifactLine(out, entry.SourceID, lockSet) } + fmt.Fprintln(out, "Extraction:") + for _, entry := range catalog.ListExtraction() { + state := "unavailable" + if entry.Available { + state = "available" + } + writeExtractionArtifactLine(out, entry.SourceID, state, entry.Provenance, lockSet) + } fmt.Fprintln(out, "Previous-session:") for _, req := range artifacts.CollectPreviousArtifactRequirements(configuredScriptoriumArtifacts(cfg)) { fmt.Fprintf(out, "- %s required=%t\n", artifactpolicy.PreviousSessionSourceID(req.Name), req.Required) @@ -50,6 +67,17 @@ func writeArtifactList(out io.Writer, cfg *config.Config, catalog *artifacts.Art } } +func writeExtractionArtifactLine(out io.Writer, source, state, provenance string, lockSet map[string]config.PublishLockRule) { + parts := []string{source, "planned", state} + if strings.TrimSpace(provenance) != "" { + parts = append(parts, "provenance="+strings.TrimSpace(provenance)) + } + if _, ok := lockSet[source]; ok { + parts = append(parts, "locked") + } + fmt.Fprintf(out, "- %s\n", strings.Join(parts, " ")) +} + func writeArtifactLine(out io.Writer, source string, lockSet map[string]config.PublishLockRule) { parts := []string{source} if _, ok := lockSet[source]; ok { @@ -103,7 +131,12 @@ func remotePublishedOutputAvailability(ctx context.Context, cfg *config.Config, func helperPublishedOutputDest(rule config.PublishOutputRule, catalog *artifacts.ArtifactCatalog) (string, bool, error) { source := strings.TrimSpace(rule.Source) - normalized, err := artifactpolicy.ResolvePublishedDestination(source, rule.Dest, helperConfiguredOutputPathMap(catalog)) + normalized, err := artifactpolicy.ResolvePublishedDestinationWithExtractions( + source, + rule.Dest, + helperConfiguredOutputPathMap(catalog), + helperExtractionOutputSet(catalog), + ) if err != nil { return "", false, err } @@ -112,6 +145,19 @@ func helperPublishedOutputDest(rule config.PublishOutputRule, catalog *artifacts return normalized, showDest, nil } +func helperExtractionOutputSet(catalog *artifacts.ArtifactCatalog) map[string]struct{} { + out := map[string]struct{}{} + if catalog == nil { + return out + } + for _, entry := range catalog.ListExtraction() { + if strings.TrimSpace(entry.ExtractionKey) != "" { + out[entry.ExtractionKey] = struct{}{} + } + } + return out +} + func helperConfiguredOutputPathMap(catalog *artifacts.ArtifactCatalog) map[string]string { out := map[string]string{} if catalog == nil { diff --git a/internal/app/operator_artifacts_list.go b/internal/app/operator_artifacts_list.go index 5960425..9ee6894 100644 --- a/internal/app/operator_artifacts_list.go +++ b/internal/app/operator_artifacts_list.go @@ -22,11 +22,11 @@ func ArtifactsList(ctx context.Context, args []string, out io.Writer) error { if strings.TrimSpace(flags.sessionID) == "" { return fmt.Errorf("artifacts list: session_id is required") } - cfg, store, locks, _, err := loadHelperContext(ctx, flags, remote) + cfg, store, locks, m, err := loadHelperContext(ctx, flags, remote) if err != nil { return fmt.Errorf("artifacts list: %w", err) } - catalog, err := buildHelperArtifactCatalog(cfg) + catalog, err := buildHelperArtifactCatalog(cfg, m) if err != nil { return fmt.Errorf("artifacts list: %w", err) } diff --git a/internal/app/operator_helpers_test.go b/internal/app/operator_helpers_test.go index a78a2d3..f77c636 100644 --- a/internal/app/operator_helpers_test.go +++ b/internal/app/operator_helpers_test.go @@ -8,8 +8,10 @@ import ( "path/filepath" "strings" "testing" + "time" "gitea.maximumdirect.net/eric/narratio/internal/adapters/storage" + "gitea.maximumdirect.net/eric/narratio/internal/artifactmodel" "gitea.maximumdirect.net/eric/narratio/internal/artifacts" "gitea.maximumdirect.net/eric/narratio/internal/config" "gitea.maximumdirect.net/eric/narratio/internal/manifest" @@ -772,6 +774,71 @@ func TestExecuteArtifactsListRemoteReportsPublishedAvailability(t *testing.T) { } } +func TestExecuteArtifactsListReportsExtractionLifecycleWithoutPayload(t *testing.T) { + workspaceRoot := t.TempDir() + pipelinePath, campaignPath, sessionPath := writeValidConfigFiles(t, workspaceRoot) + addExtractionOutputToPipeline(t, pipelinePath) + addPublishOutputsToPipeline(t, pipelinePath, ` + outputs: + - source: narratio.extraction.encounters + dest: artifacts/encounters.json + required: true +`) + lanePath := writeOperatorExtractionManifest(t, workspaceRoot) + fake := &storage.FakeBackend{} + publishedKey := artifacts.S3PublishedOutputKey( + artifacts.S3SessionPrefix("dnd", "sample-campaign", "2026-05-03"), + "artifacts/encounters.json", + ) + fake.SeedObject(storage.FakeObject{Key: publishedKey, Data: []byte(`{"secret":"DO_NOT_PRINT"}`)}) + var storeInitCalls int + restoreAppConfigTestGlobals(t, fake, &storeInitCalls, []string{sessionPath}) + + var stdout bytes.Buffer + var stderr bytes.Buffer + code := Execute([]string{ + "session", "artifacts", "2026-05-03", + "--config", pipelinePath, + "--campaign-file", campaignPath, + "--session", sessionPath, + "--remote", + }, &stdout, &stderr) + if code != 0 { + t.Fatalf("exit code = %d, want 0; stderr=%q", code, stderr.String()) + } + out := stdout.String() + for _, want := range []string{ + "Extraction:", + "narratio.extraction.encounters planned available provenance=manifest.current_extract_run", + "narratio.extraction.encounters dest=artifacts/encounters.json remote=published", + } { + if !strings.Contains(out, want) { + t.Fatalf("stdout = %q, want %q", out, want) + } + } + if strings.Contains(out, "DO_NOT_PRINT") { + t.Fatalf("operator output exposed extraction payload: %q", out) + } + + if err := os.Remove(lanePath); err != nil { + t.Fatal(err) + } + stdout.Reset() + stderr.Reset() + code = Execute([]string{ + "session", "artifacts", "2026-05-03", + "--config", pipelinePath, + "--campaign-file", campaignPath, + "--session", sessionPath, + }, &stdout, &stderr) + if code != 0 { + t.Fatalf("unavailable exit code = %d, want 0; stderr=%q", code, stderr.String()) + } + if !strings.Contains(stdout.String(), "narratio.extraction.encounters planned unavailable") { + t.Fatalf("stdout = %q, want unavailable extraction state", stdout.String()) + } +} + func TestExecuteArtifactsListRemoteUsesPublishOutputDestinations(t *testing.T) { workspaceRoot := t.TempDir() pipelinePath, campaignPath, sessionPath := writeValidConfigFiles(t, workspaceRoot) @@ -1041,6 +1108,72 @@ func addPublishOutputsToPipeline(t *testing.T, pipelinePath, publishYAML string) } } +func addExtractionOutputToPipeline(t *testing.T, pipelinePath string) { + t.Helper() + data, err := os.ReadFile(pipelinePath) + if err != nil { + t.Fatal(err) + } + data = append(data, []byte(`notarius: + enabled: true + config_path: notarius.yml + pipeline_id: campaign.extract + outputs: + encounters: + lane_id: encounters + media_type: application/json + schema_id: encounters + schema_version: "1" + module_key: encounters +`)...) + if err := os.WriteFile(pipelinePath, data, 0o644); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(filepath.Join(filepath.Dir(pipelinePath), "notarius.yml"), []byte("{}\n"), 0o644); err != nil { + t.Fatal(err) + } +} + +func writeOperatorExtractionManifest(t *testing.T, workspaceRoot string) string { + t.Helper() + paths := artifacts.NewLocalStore(workspaceRoot).SessionPathsFor("sample-campaign", "2026-05-03") + bundleRoot := filepath.Join(paths.ArtifactsDir, "notarius", "extract-run-1") + lanePath := filepath.Join(bundleRoot, "lanes", "encounters.json") + indexPath := filepath.Join(bundleRoot, "index.json") + mustWriteTestFile(t, lanePath, `{"secret":"DO_NOT_PRINT"}`) + mustWriteTestFile(t, indexPath, `{"lanes":[]}`) + laneChecksum, err := artifacts.SHA256File(lanePath) + if err != nil { + t.Fatal(err) + } + indexChecksum, err := artifacts.SHA256File(indexPath) + if err != nil { + t.Fatal(err) + } + m := manifest.New("2026-05-03", time.Now().UTC()) + m.Campaign = "sample-campaign" + m.Stages["extract"] = &manifest.StageRecord{ + Name: "extract", Status: manifest.StatusSucceeded, + Metadata: map[string]any{ + "narratio_run_id": "extract-run-1", "bundle_root": bundleRoot, + "receipt": map[string]any{"run_id": "notarius-run-1", "pipeline_id": "campaign.extract"}, + }, + Outputs: []manifest.ArtifactRecord{ + { + Kind: "notarius_lane", SourceID: artifacts.ExtractionArtifactSourceID("encounters"), LocalPath: lanePath, + ProducerRunID: "extract-run-1", Checksum: laneChecksum, + Contract: &artifactmodel.ContractMetadata{MediaType: "application/json", SchemaID: "encounters", SchemaVersion: "1", ModuleKey: "encounters"}, + ExternalProvenance: &artifactmodel.ExternalProvenance{System: "notarius", RunID: "notarius-run-1", PipelineID: "campaign.extract", ArtifactID: "encounters"}, + }, + {Kind: "notarius_index", LocalPath: indexPath, ProducerRunID: "extract-run-1", Checksum: indexChecksum}, + }, + } + if err := (&manifest.LocalStore{}).Save(context.Background(), paths.ManifestPath, m); err != nil { + t.Fatal(err) + } + return lanePath +} + func replaceInFileOrFatal(t *testing.T, path, old, new string) { t.Helper() data, err := os.ReadFile(path) diff --git a/internal/app/operator_status.go b/internal/app/operator_status.go index eb9ce06..704137b 100644 --- a/internal/app/operator_status.go +++ b/internal/app/operator_status.go @@ -11,6 +11,7 @@ import ( "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" ) // Status reports effective local/remote session state. @@ -40,11 +41,13 @@ func Status(ctx context.Context, args []string, out io.Writer) error { writeStatusStableInputs(out, inspectStableInputs(cfg)) writeStatusLocalAudio(out, inspectLocalAudioPresence(cfg)) + var localManifest *manifest.Manifest if m, err := loadLocalManifest(ctx, paths.ManifestPath); err != nil { fmt.Fprintf(out, "Local manifest: error: %v\n", err) } else if m == nil { fmt.Fprintln(out, "Local manifest: missing") } else { + localManifest = m fmt.Fprintf(out, "Local manifest: %s\n", paths.ManifestPath) writeStageStatuses(out, m) } @@ -72,7 +75,7 @@ func Status(ctx context.Context, args []string, out io.Writer) error { lockChecks := inspectEffectiveLocks(ctx, cfg, store) locks := lockChecks.Locks lockErr := lockChecks.Err - if catalog, catalogErr := buildHelperArtifactCatalog(cfg); catalogErr != nil { + if catalog, catalogErr := buildHelperArtifactCatalog(cfg, localManifest); catalogErr != nil { fmt.Fprintf(out, "Remote outputs: error: %v\n", catalogErr) } else if storeErr == nil { catalogLocks := locks diff --git a/internal/app/restore_execution_test.go b/internal/app/restore_execution_test.go index 1cb3e6e..225aca9 100644 --- a/internal/app/restore_execution_test.go +++ b/internal/app/restore_execution_test.go @@ -9,8 +9,10 @@ import ( "path/filepath" "strings" "testing" + "time" "gitea.maximumdirect.net/eric/narratio/internal/adapters/storage" + "gitea.maximumdirect.net/eric/narratio/internal/artifactmodel" "gitea.maximumdirect.net/eric/narratio/internal/artifacts" "gitea.maximumdirect.net/eric/narratio/internal/config" "gitea.maximumdirect.net/eric/narratio/internal/manifest" @@ -62,6 +64,55 @@ func TestExecuteRestoreNonDryRunRestoresDurableFiles(t *testing.T) { } } +func TestExecuteRestoreRoundTripsPublishedExtractionAndManifestMetadata(t *testing.T) { + workspaceRoot := t.TempDir() + pipelinePath, campaignPath, sessionPath := writeValidConfigFiles(t, workspaceRoot) + fake := &storage.FakeBackend{} + cfg, sessionPrefix, manifestKey, runIDKey := seedRestoreCommittedState(t, fake, pipelinePath, campaignPath, sessionPath) + seedRestoreObject(fake, sessionPrefix+"artifacts/encounters.json", []byte(`{"encounters":[]}`)) + seedRestoreObject(fake, runIDKey, []byte("20260519T010203Z-a1b2c3d4\n")) + + remoteManifest := manifest.New(cfg.Session.SessionID, time.Now().UTC()) + remoteManifest.Campaign = cfg.Session.Campaign + remoteManifest.RunID = "20260519T010203Z-a1b2c3d4" + remoteManifest.Stages["extract"] = &manifest.StageRecord{ + Name: "extract", Status: manifest.StatusSucceeded, + Outputs: []manifest.ArtifactRecord{{ + Kind: "notarius_lane", SourceID: artifacts.ExtractionArtifactSourceID("encounters"), + LocalPath: "/prior/workspace/artifacts/notarius/extract-run-1/lanes/encounters.json", + Contract: &artifactmodel.ContractMetadata{ + MediaType: "application/json", SchemaID: "encounters", SchemaVersion: "1", + }, + ExternalProvenance: &artifactmodel.ExternalProvenance{ + System: "notarius", RunID: "notarius-run-1", PipelineID: "campaign.extract", ArtifactID: "encounters", + }, + }}, + } + manifestBody, err := json.Marshal(remoteManifest) + if err != nil { + t.Fatal(err) + } + seedRestoreObject(fake, manifestKey, manifestBody) + restoreWithStoreAndRealPhases(t, fake) + + var stdout bytes.Buffer + var stderr bytes.Buffer + code := Execute([]string{"session", "restore", "2026-05-03", "--config", pipelinePath, "--campaign-file", campaignPath, "--session", sessionPath}, &stdout, &stderr) + if code != 0 { + t.Fatalf("exit code = %d, want 0; stderr=%q", code, stderr.String()) + } + sessionRoot := artifacts.SessionWorkDirForCampaign(workspaceRoot, cfg.Session.Campaign, cfg.Session.SessionID) + mustReadEquals(t, filepath.Join(sessionRoot, "artifacts", "encounters.json"), `{"encounters":[]}`) + restored, err := (&manifest.LocalStore{}).Load(context.Background(), filepath.Join(sessionRoot, "manifest.json")) + if err != nil { + t.Fatalf("load restored manifest: %v", err) + } + lane := restored.Stages["extract"].Outputs[0] + if lane.Contract == nil || lane.Contract.SchemaID != "encounters" || lane.ExternalProvenance == nil || lane.ExternalProvenance.RunID != "notarius-run-1" { + t.Fatalf("restored extraction metadata = %#v", lane) + } +} + func TestExecuteRestoreIncludeAudioRestoresAudio(t *testing.T) { workspaceRoot := t.TempDir() pipelinePath, campaignPath, sessionPath := writeValidConfigFiles(t, workspaceRoot) diff --git a/internal/artifactpolicy/policy.go b/internal/artifactpolicy/policy.go index 17c64a8..173b7ed 100644 --- a/internal/artifactpolicy/policy.go +++ b/internal/artifactpolicy/policy.go @@ -311,7 +311,17 @@ func DeriveDefaultPublishedDestination(source Source, configured map[string]stri // ResolvePublishedDestination validates and normalizes an explicit destination, // or derives one when omitted. func ResolvePublishedDestination(sourceID, explicitDest string, configured map[string]string) (string, error) { - source, err := ValidatePublishSource(sourceID, configured) + return ResolvePublishedDestinationWithExtractions(sourceID, explicitDest, configured, nil) +} + +// ResolvePublishedDestinationWithExtractions validates and normalizes a destination +// while accepting extraction sources declared by the effective Notarius configuration. +func ResolvePublishedDestinationWithExtractions( + sourceID, explicitDest string, + configured map[string]string, + extractions map[string]struct{}, +) (string, error) { + source, err := ValidatePublishSourceWithExtractions(sourceID, configured, extractions) if err != nil { return "", err } diff --git a/internal/artifacts/catalog.go b/internal/artifacts/catalog.go index d1bbf94..5e4a7f3 100644 --- a/internal/artifacts/catalog.go +++ b/internal/artifacts/catalog.go @@ -6,6 +6,7 @@ import ( "strings" "gitea.maximumdirect.net/eric/narratio/internal/artifactpolicy" + "gitea.maximumdirect.net/eric/narratio/internal/config" ) const ( @@ -30,6 +31,25 @@ type ExtractionArtifactDefinition struct { ModuleKey string } +// ExtractionDefinitionsFromConfig converts the effective Notarius output map into catalog definitions. +func ExtractionDefinitionsFromConfig(cfg *config.NotariusConfig) map[string]ExtractionArtifactDefinition { + if cfg == nil || len(cfg.Outputs) == 0 { + return nil + } + definitions := make(map[string]ExtractionArtifactDefinition, len(cfg.Outputs)) + for key, output := range cfg.Outputs { + definitions[key] = ExtractionArtifactDefinition{ + LaneID: output.LaneID, + PipelineID: cfg.PipelineID, + MediaType: output.MediaType, + SchemaID: output.SchemaID, + SchemaVersion: output.SchemaVersion, + ModuleKey: output.ModuleKey, + } + } + return definitions +} + // CatalogEntry is one runtime catalog entry resolved by source ID. type CatalogEntry struct { SourceID string diff --git a/internal/config/validate.go b/internal/config/validate.go index 4845d6f..5c21380 100644 --- a/internal/config/validate.go +++ b/internal/config/validate.go @@ -172,7 +172,7 @@ func validatePublish(cfg *PublishConfig, scriptorium *ScriptoriumConfig, notariu } dest := strings.TrimSpace(item.Dest) if dest == "" { - derivedDest, err := artifactpolicy.ResolvePublishedDestination(source, "", configuredOutputs) + derivedDest, err := artifactpolicy.ResolvePublishedDestinationWithExtractions(source, "", configuredOutputs, extractionOutputs) if err != nil { return fmt.Errorf("%s.dest is required when destination cannot be derived from %q: %w", prefix, source, err) } diff --git a/internal/stage/analyze.go b/internal/stage/analyze.go index 2bc4378..9d34300 100644 --- a/internal/stage/analyze.go +++ b/internal/stage/analyze.go @@ -752,11 +752,13 @@ func buildAnalyzeRuntimeArtifactCatalog( if err := catalog.RegisterBuiltIns(); err != nil { return nil, err } - extractionDefinitions := configuredExtractionDefinitions(notariusCfg) + extractionDefinitions := artifacts.ExtractionDefinitionsFromConfig(notariusCfg) if err := catalog.RegisterExtractionArtifacts(extractionDefinitions); err != nil { return nil, err } - catalog.HydrateExtractionArtifacts(paths, m, extractionDefinitions) + if notariusCfg != nil && notariusCfg.Enabled { + catalog.HydrateExtractionArtifacts(paths, m, extractionDefinitions) + } if scriptoriumCfg == nil { return catalog, nil } diff --git a/internal/stage/publish.go b/internal/stage/publish.go index 3820042..68e8d51 100644 --- a/internal/stage/publish.go +++ b/internal/stage/publish.go @@ -124,7 +124,12 @@ func (publishStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*S if err != nil { return nil, fmt.Errorf("publish: collect previous files: %w", err) } - runtimeCatalog, err := buildPublishRuntimeArtifactCatalog(sessionPaths, env.Config.Pipeline.Scriptorium) + runtimeCatalog, err := buildPublishRuntimeArtifactCatalog( + sessionPaths, + m, + env.Config.Pipeline.Scriptorium, + env.Config.Pipeline.Notarius, + ) if err != nil { return nil, fmt.Errorf("publish: build runtime artifact catalog: %w", err) } @@ -361,10 +366,11 @@ func resolvePublishOutputs( lockSet := publishLockSet(locks) selectedSet := publishSelectedArtifactSet(selectedArtifactKeys) configuredOutputs := configuredOutputPathMapFromCatalog(catalog) + extractionOutputs := extractionOutputSetFromCatalog(catalog) for _, rule := range rules { source := strings.TrimSpace(rule.Source) required := rule.Required == nil || *rule.Required - dest, err := resolvePublishOutputDest(rule, configuredOutputs) + dest, err := resolvePublishOutputDest(rule, configuredOutputs, extractionOutputs) if err != nil { return nil, nil, nil, nil, fmt.Errorf("source %q: %w", source, err) } @@ -454,8 +460,8 @@ func publishLockSet(locks []config.PublishLockRule) map[string]config.PublishLoc return out } -func resolvePublishOutputDest(rule config.PublishOutputRule, configured map[string]string) (string, error) { - return artifactpolicy.ResolvePublishedDestination(rule.Source, rule.Dest, configured) +func resolvePublishOutputDest(rule config.PublishOutputRule, configured map[string]string, extractions map[string]struct{}) (string, error) { + return artifactpolicy.ResolvePublishedDestinationWithExtractions(rule.Source, rule.Dest, configured, extractions) } func configuredOutputPathMapFromCatalog(catalog *artifacts.ArtifactCatalog) map[string]string { @@ -472,14 +478,36 @@ func configuredOutputPathMapFromCatalog(catalog *artifacts.ArtifactCatalog) map[ return out } +func extractionOutputSetFromCatalog(catalog *artifacts.ArtifactCatalog) map[string]struct{} { + out := map[string]struct{}{} + if catalog == nil { + return out + } + for _, entry := range catalog.ListExtraction() { + if strings.TrimSpace(entry.ExtractionKey) != "" { + out[entry.ExtractionKey] = struct{}{} + } + } + return out +} + func buildPublishRuntimeArtifactCatalog( paths artifacts.SessionPaths, + m *manifest.Manifest, scriptoriumCfg *config.ScriptoriumConfig, + notariusCfg *config.NotariusConfig, ) (*artifacts.ArtifactCatalog, error) { catalog := artifacts.NewArtifactCatalog() if err := catalog.RegisterBuiltIns(); err != nil { return nil, err } + extractionDefinitions := artifacts.ExtractionDefinitionsFromConfig(notariusCfg) + if err := catalog.RegisterExtractionArtifacts(extractionDefinitions); err != nil { + return nil, err + } + if notariusCfg != nil && notariusCfg.Enabled { + catalog.HydrateExtractionArtifacts(paths, m, extractionDefinitions) + } if scriptoriumCfg == nil { return catalog, nil } diff --git a/internal/stage/publish_test.go b/internal/stage/publish_test.go index 2d3b449..ce8442c 100644 --- a/internal/stage/publish_test.go +++ b/internal/stage/publish_test.go @@ -12,6 +12,7 @@ import ( "time" "gitea.maximumdirect.net/eric/narratio/internal/adapters/storage" + "gitea.maximumdirect.net/eric/narratio/internal/artifactmodel" "gitea.maximumdirect.net/eric/narratio/internal/artifacts" "gitea.maximumdirect.net/eric/narratio/internal/config" "gitea.maximumdirect.net/eric/narratio/internal/manifest" @@ -198,6 +199,119 @@ func TestPublishUsesCustomOutputRules(t *testing.T) { } } +func TestPublishUploadsExplicitExtractionAndPreservesManifestMetadata(t *testing.T) { + env, m, _ := publishFixture(t) + lanePath := configurePublishExtractionFixture(t, env, m) + env.SelectedArtifactKeys = []string{"session_recap"} + env.Config.Pipeline.Publish.Outputs = []config.PublishOutputRule{ + {Source: artifacts.ExtractionArtifactSourceID("encounters"), Dest: "artifacts/encounters.json", Required: boolPtr(true)}, + } + + result, err := publishStage{}.Run(context.Background(), env, m) + if err != nil { + t.Fatalf("Run() error = %v", err) + } + fake := env.ObjectStore.(*storage.FakeBackend) + publishedKey := m.S3SessionPrefix + "artifacts/encounters.json" + if got := string(fake.Objects[publishedKey].Data); got != `{"encounters":[]}` { + t.Fatalf("published extraction = %q", got) + } + if _, ok := fake.Objects[m.S3SessionPrefix+"artifacts/notarius/index.json"]; ok { + t.Fatal("Notarius index was implicitly published") + } + for key := range fake.Objects { + if strings.Contains(key, "/notarius/") { + t.Fatalf("Notarius bundle member was implicitly published at %q", key) + } + } + if result.Metadata["published_files_uploaded"] != 1 { + t.Fatalf("published_files_uploaded = %#v, want 1", result.Metadata["published_files_uploaded"]) + } + if uploads := fake.Uploads; len(uploads) == 0 || uploads[len(uploads)-1].Key != m.S3SessionPrefix+"current/run_id.txt" { + t.Fatalf("last upload = %#v, want current run pointer", uploads) + } + + var current manifest.Manifest + if err := json.Unmarshal(fake.Objects[m.S3SessionPrefix+"current/manifest.json"].Data, ¤t); err != nil { + t.Fatalf("unmarshal current manifest: %v", err) + } + extractRecord := current.Stages["extract"] + if extractRecord == nil || len(extractRecord.Outputs) != 2 { + t.Fatalf("extract record = %#v", extractRecord) + } + lane := extractRecord.Outputs[0] + if lane.LocalPath != lanePath || lane.Contract == nil || lane.Contract.SchemaID != "encounters" || + lane.ExternalProvenance == nil || lane.ExternalProvenance.RunID != "notarius-run-1" { + t.Fatalf("serialized extraction metadata = %#v", lane) + } +} + +func TestPublishRequiredInvalidExtractionFails(t *testing.T) { + env, m, _ := publishFixture(t) + lanePath := configurePublishExtractionFixture(t, env, m) + env.Config.Pipeline.Publish.Outputs = []config.PublishOutputRule{ + {Source: artifacts.ExtractionArtifactSourceID("encounters"), Dest: "artifacts/encounters.json", Required: boolPtr(true)}, + } + writeStageTestFile(t, lanePath, `{"tampered":true}`) + + _, err := publishStage{}.Run(context.Background(), env, m) + if err == nil || !strings.Contains(err.Error(), `required output source unavailable: "narratio.extraction.encounters"`) { + t.Fatalf("Run() error = %v, want unavailable extraction failure", err) + } + if len(env.ObjectStore.(*storage.FakeBackend).Uploads) != 0 { + t.Fatal("publish uploaded files after extraction validation failed") + } +} + +func TestPublishDisabledExtractionConfigurationIsUnavailable(t *testing.T) { + env, m, _ := publishFixture(t) + configurePublishExtractionFixture(t, env, m) + env.Config.Pipeline.Notarius.Enabled = false + env.Config.Pipeline.Publish.Outputs = []config.PublishOutputRule{ + {Source: artifacts.ExtractionArtifactSourceID("encounters"), Dest: "artifacts/encounters.json", Required: boolPtr(true)}, + } + + _, err := publishStage{}.Run(context.Background(), env, m) + if err == nil || !strings.Contains(err.Error(), `required output source unavailable: "narratio.extraction.encounters"`) { + t.Fatalf("Run() error = %v, want unavailable disabled extraction", err) + } +} + +func TestPublishOptionalMissingExtractionIsSkipped(t *testing.T) { + env, m, _ := publishFixture(t) + lanePath := configurePublishExtractionFixture(t, env, m) + env.Config.Pipeline.Publish.Outputs = []config.PublishOutputRule{ + {Source: artifacts.ExtractionArtifactSourceID("encounters"), Dest: "artifacts/encounters.json", Required: boolPtr(false)}, + } + if err := os.Remove(lanePath); err != nil { + t.Fatal(err) + } + + result, err := publishStage{}.Run(context.Background(), env, m) + if err != nil { + t.Fatalf("Run() error = %v", err) + } + if got := result.Metadata["skipped_optional_outputs"].([]string); !reflect.DeepEqual(got, []string{"artifacts/encounters.json"}) { + t.Fatalf("skipped_optional_outputs = %#v", got) + } +} + +func TestPublishArtifactSelectionDoesNotFilterExtractionOutputs(t *testing.T) { + env, m, _ := publishFixture(t) + configurePublishExtractionFixture(t, env, m) + env.SelectedArtifactKeys = []string{"session_recap"} + env.Config.Pipeline.Publish.Outputs = []config.PublishOutputRule{ + {Source: artifacts.ExtractionArtifactSourceID("encounters"), Dest: "artifacts/encounters.json", Required: boolPtr(true)}, + } + + if _, err := (publishStage{}).Run(context.Background(), env, m); err != nil { + t.Fatalf("Run() error = %v", err) + } + if _, ok := env.ObjectStore.(*storage.FakeBackend).Objects[m.S3SessionPrefix+"artifacts/encounters.json"]; !ok { + t.Fatal("explicit extraction output was filtered by Scriptorium artifact selection") + } +} + func TestPublishSkipsOptionalMissingOutput(t *testing.T) { env, m, _ := publishFixture(t) env.Config.Pipeline.Publish.Outputs = []config.PublishOutputRule{ @@ -591,6 +705,51 @@ func publishFixture(t *testing.T) (*Env, *manifest.Manifest, string) { return env, m, runRoot } +func configurePublishExtractionFixture(t *testing.T, env *Env, m *manifest.Manifest) string { + t.Helper() + env.Config.Pipeline.Notarius = &config.NotariusConfig{ + Enabled: true, PipelineID: "campaign.extract", + Outputs: map[string]config.NotariusOutputConfig{ + "encounters": { + LaneID: "encounters", MediaType: "application/json", SchemaID: "encounters", + SchemaVersion: "1", ModuleKey: "encounters", + }, + }, + } + paths := publishSessionPaths(env, m) + producerRunID := "extract-run-1" + bundleRoot := filepath.Join(paths.ArtifactsDir, "notarius", producerRunID) + lanePath := filepath.Join(bundleRoot, "lanes", "encounters.json") + indexPath := filepath.Join(bundleRoot, "index.json") + writeStageTestFile(t, lanePath, `{"encounters":[]}`) + writeStageTestFile(t, indexPath, `{"lanes":[]}`) + laneChecksum, err := artifacts.SHA256File(lanePath) + if err != nil { + t.Fatal(err) + } + indexChecksum, err := artifacts.SHA256File(indexPath) + if err != nil { + t.Fatal(err) + } + m.Stages["extract"] = &manifest.StageRecord{ + Name: "extract", Status: manifest.StatusSucceeded, + Metadata: map[string]any{ + "narratio_run_id": producerRunID, "bundle_root": bundleRoot, + "receipt": map[string]any{"run_id": "notarius-run-1", "pipeline_id": "campaign.extract"}, + }, + Outputs: []manifest.ArtifactRecord{ + { + Kind: "notarius_lane", SourceID: artifacts.ExtractionArtifactSourceID("encounters"), LocalPath: lanePath, + ProducerRunID: producerRunID, Checksum: laneChecksum, + Contract: &artifactmodel.ContractMetadata{MediaType: "application/json", SchemaID: "encounters", SchemaVersion: "1", ModuleKey: "encounters"}, + ExternalProvenance: &artifactmodel.ExternalProvenance{System: "notarius", RunID: "notarius-run-1", PipelineID: "campaign.extract", ArtifactID: "encounters"}, + }, + {Kind: "notarius_index", LocalPath: indexPath, ProducerRunID: producerRunID, Checksum: indexChecksum}, + }, + } + return lanePath +} + type publishedOutputFailingStore struct { delegate *storage.FakeBackend failKey string diff --git a/internal/stage/runtime_catalog.go b/internal/stage/runtime_catalog.go deleted file mode 100644 index 5543162..0000000 --- a/internal/stage/runtime_catalog.go +++ /dev/null @@ -1,24 +0,0 @@ -package stage - -import ( - "gitea.maximumdirect.net/eric/narratio/internal/artifacts" - "gitea.maximumdirect.net/eric/narratio/internal/config" -) - -func configuredExtractionDefinitions(cfg *config.NotariusConfig) map[string]artifacts.ExtractionArtifactDefinition { - if cfg == nil || !cfg.Enabled || len(cfg.Outputs) == 0 { - return nil - } - definitions := make(map[string]artifacts.ExtractionArtifactDefinition, len(cfg.Outputs)) - for key, output := range cfg.Outputs { - definitions[key] = artifacts.ExtractionArtifactDefinition{ - LaneID: output.LaneID, - PipelineID: cfg.PipelineID, - MediaType: output.MediaType, - SchemaID: output.SchemaID, - SchemaVersion: output.SchemaVersion, - ModuleKey: output.ModuleKey, - } - } - return definitions -}