Publish and inspect configured extraction artifacts
This commit is contained in:
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
}
|
||||
Reference in New Issue
Block a user