Skip to main content

dag_ml_core/
provenance.rs

1use std::collections::{BTreeMap, BTreeSet};
2
3use serde::{de::DeserializeOwned, Deserialize, Serialize};
4use serde_json::{json, Value};
5use sha2::{Digest, Sha256};
6
7use crate::bundle::ExecutionBundle;
8use crate::data::ExternalDataPlanEnvelope;
9use crate::error::{DagMlError, Result};
10use crate::ids::{ArtifactId, LineageId};
11use crate::plan::ExecutionPlan;
12use crate::runtime::{
13    FileArtifactManifest, FilePredictionCacheManifest, LineageRecord, FILE_ARTIFACT_MANIFEST_FILE,
14    FILE_PREDICTION_CACHE_MANIFEST_FILE,
15};
16
17pub const RESEARCH_PROVENANCE_SCHEMA_VERSION: u32 = 1;
18pub const EXECUTION_PLAN_FILE: &str = "execution_plan.json";
19pub const EXECUTION_BUNDLE_FILE: &str = "execution_bundle.json";
20pub const LINEAGE_RECORDS_FILE: &str = "lineage_records.json";
21pub const PROV_JSONLD_FILE: &str = "lineage.prov.jsonld";
22pub const RO_CRATE_METADATA_FILE: &str = "ro-crate-metadata.json";
23pub const OPENLINEAGE_RUN_EVENT_SCHEMA_URL: &str =
24    "https://openlineage.io/spec/1-0-0/OpenLineage.json#/definitions/RunEvent";
25pub const DAGML_OPENLINEAGE_FACET_SCHEMA_URL: &str =
26    "https://github.com/GBeurier/dag-ml/schemas/openlineage_dagml_facets.v1.schema.json";
27
28#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
29pub struct ResearchProvenanceExport {
30    pub schema_version: u32,
31    pub prov_jsonld: Value,
32    pub ro_crate_metadata: Value,
33}
34
35#[derive(Clone, Debug, Eq, PartialEq)]
36pub struct ResearchProvenancePackage {
37    pub schema_version: u32,
38    pub files: BTreeMap<String, ResearchProvenancePackageFile>,
39}
40
41#[derive(Clone, Debug, Eq, PartialEq)]
42pub struct ResearchProvenancePackageFile {
43    pub path: String,
44    pub sha256: String,
45    pub size_bytes: usize,
46    pub bytes: Vec<u8>,
47}
48
49#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
50pub struct ResearchProvenancePackageValidation {
51    pub schema_version: u32,
52    pub plan_id: String,
53    pub bundle_id: String,
54    pub file_count: usize,
55    pub checksummed_file_count: usize,
56    pub lineage_record_count: usize,
57    pub data_envelope_count: usize,
58    pub has_prediction_cache_manifest: bool,
59    pub has_artifact_manifest: bool,
60}
61
62#[derive(Clone, Debug, Eq, PartialEq)]
63pub struct OpenLineageRunEventOptions {
64    pub namespace: String,
65    pub event_time: String,
66}
67
68pub fn build_research_provenance_export(
69    plan: &ExecutionPlan,
70    bundle: &ExecutionBundle,
71    lineage: &[LineageRecord],
72    data_envelopes: &BTreeMap<String, ExternalDataPlanEnvelope>,
73    prediction_cache_manifest: Option<&FilePredictionCacheManifest>,
74    artifact_manifest: Option<&FileArtifactManifest>,
75) -> Result<ResearchProvenanceExport> {
76    validate_provenance_inputs(
77        plan,
78        bundle,
79        lineage,
80        data_envelopes,
81        prediction_cache_manifest,
82        artifact_manifest,
83    )?;
84
85    Ok(ResearchProvenanceExport {
86        schema_version: RESEARCH_PROVENANCE_SCHEMA_VERSION,
87        prov_jsonld: build_prov_jsonld(
88            plan,
89            bundle,
90            lineage,
91            data_envelopes,
92            prediction_cache_manifest,
93            artifact_manifest,
94        )?,
95        ro_crate_metadata: build_ro_crate_metadata(
96            plan,
97            bundle,
98            data_envelopes,
99            prediction_cache_manifest,
100            artifact_manifest,
101        )?,
102    })
103}
104
105pub fn build_research_provenance_package(
106    plan: &ExecutionPlan,
107    bundle: &ExecutionBundle,
108    lineage: &[LineageRecord],
109    data_envelopes: &BTreeMap<String, ExternalDataPlanEnvelope>,
110    prediction_cache_manifest: Option<&FilePredictionCacheManifest>,
111    artifact_manifest: Option<&FileArtifactManifest>,
112) -> Result<ResearchProvenancePackage> {
113    let export = build_research_provenance_export(
114        plan,
115        bundle,
116        lineage,
117        data_envelopes,
118        prediction_cache_manifest,
119        artifact_manifest,
120    )?;
121    let mut files = BTreeMap::new();
122    add_json_package_file(&mut files, EXECUTION_PLAN_FILE, plan, "execution plan")?;
123    add_json_package_file(
124        &mut files,
125        EXECUTION_BUNDLE_FILE,
126        bundle,
127        "execution bundle",
128    )?;
129    add_json_package_file(
130        &mut files,
131        LINEAGE_RECORDS_FILE,
132        &lineage,
133        "lineage records",
134    )?;
135    add_json_package_file(
136        &mut files,
137        PROV_JSONLD_FILE,
138        &export.prov_jsonld,
139        "PROV JSON-LD",
140    )?;
141    if let Some(manifest) = prediction_cache_manifest {
142        add_json_package_file(
143            &mut files,
144            FILE_PREDICTION_CACHE_MANIFEST_FILE,
145            manifest,
146            "prediction cache manifest",
147        )?;
148    }
149    if let Some(manifest) = artifact_manifest {
150        add_json_package_file(
151            &mut files,
152            FILE_ARTIFACT_MANIFEST_FILE,
153            manifest,
154            "artifact manifest",
155        )?;
156    }
157    for (key, envelope) in data_envelopes {
158        add_json_package_file(
159            &mut files,
160            &data_envelope_file_path(key)?,
161            envelope,
162            "data envelope",
163        )?;
164    }
165
166    let mut ro_crate_metadata = export.ro_crate_metadata;
167    annotate_ro_crate_package_files(&mut ro_crate_metadata, &files)?;
168    add_json_package_file(
169        &mut files,
170        RO_CRATE_METADATA_FILE,
171        &ro_crate_metadata,
172        "RO-Crate metadata",
173    )?;
174
175    Ok(ResearchProvenancePackage {
176        schema_version: RESEARCH_PROVENANCE_SCHEMA_VERSION,
177        files,
178    })
179}
180
181pub fn validate_research_provenance_package_files(
182    files: &BTreeMap<String, Vec<u8>>,
183) -> Result<ResearchProvenancePackageValidation> {
184    if files.is_empty() {
185        return Err(DagMlError::RuntimeValidation(
186            "research provenance package has no files".to_string(),
187        ));
188    }
189    for path in files.keys() {
190        validate_package_path(path)?;
191    }
192    require_package_file(files, EXECUTION_PLAN_FILE)?;
193    require_package_file(files, EXECUTION_BUNDLE_FILE)?;
194    require_package_file(files, LINEAGE_RECORDS_FILE)?;
195    require_package_file(files, PROV_JSONLD_FILE)?;
196    let ro_crate_metadata: Value = parse_package_json(
197        require_package_file(files, RO_CRATE_METADATA_FILE)?,
198        RO_CRATE_METADATA_FILE,
199    )?;
200
201    let checksummed_file_count = validate_ro_crate_package_checksums(&ro_crate_metadata, files)?;
202    validate_prov_jsonld_root(parse_package_json(
203        require_package_file(files, PROV_JSONLD_FILE)?,
204        PROV_JSONLD_FILE,
205    )?)?;
206
207    let plan: ExecutionPlan = parse_package_json(
208        require_package_file(files, EXECUTION_PLAN_FILE)?,
209        EXECUTION_PLAN_FILE,
210    )?;
211    let bundle: ExecutionBundle = parse_package_json(
212        require_package_file(files, EXECUTION_BUNDLE_FILE)?,
213        EXECUTION_BUNDLE_FILE,
214    )?;
215    let lineage: Vec<LineageRecord> = parse_package_json(
216        require_package_file(files, LINEAGE_RECORDS_FILE)?,
217        LINEAGE_RECORDS_FILE,
218    )?;
219    let data_envelopes = parse_package_data_envelopes(files)?;
220    let prediction_cache_manifest: Option<FilePredictionCacheManifest> = files
221        .get(FILE_PREDICTION_CACHE_MANIFEST_FILE)
222        .map(|bytes| parse_package_json(bytes, FILE_PREDICTION_CACHE_MANIFEST_FILE))
223        .transpose()?;
224    let artifact_manifest: Option<FileArtifactManifest> = files
225        .get(FILE_ARTIFACT_MANIFEST_FILE)
226        .map(|bytes| parse_package_json(bytes, FILE_ARTIFACT_MANIFEST_FILE))
227        .transpose()?;
228
229    validate_provenance_inputs(
230        &plan,
231        &bundle,
232        &lineage,
233        &data_envelopes,
234        prediction_cache_manifest.as_ref(),
235        artifact_manifest.as_ref(),
236    )?;
237
238    Ok(ResearchProvenancePackageValidation {
239        schema_version: RESEARCH_PROVENANCE_SCHEMA_VERSION,
240        plan_id: plan.id.to_string(),
241        bundle_id: bundle.bundle_id.to_string(),
242        file_count: files.len(),
243        checksummed_file_count,
244        lineage_record_count: lineage.len(),
245        data_envelope_count: data_envelopes.len(),
246        has_prediction_cache_manifest: prediction_cache_manifest.is_some(),
247        has_artifact_manifest: artifact_manifest.is_some(),
248    })
249}
250
251pub fn build_openlineage_run_event_from_package_files(
252    files: &BTreeMap<String, Vec<u8>>,
253    namespace: &str,
254    event_time: &str,
255) -> Result<Value> {
256    validate_research_provenance_package_files(files)?;
257    let plan: ExecutionPlan = parse_package_json(
258        require_package_file(files, EXECUTION_PLAN_FILE)?,
259        EXECUTION_PLAN_FILE,
260    )?;
261    let bundle: ExecutionBundle = parse_package_json(
262        require_package_file(files, EXECUTION_BUNDLE_FILE)?,
263        EXECUTION_BUNDLE_FILE,
264    )?;
265    let lineage: Vec<LineageRecord> = parse_package_json(
266        require_package_file(files, LINEAGE_RECORDS_FILE)?,
267        LINEAGE_RECORDS_FILE,
268    )?;
269    let data_envelopes = parse_package_data_envelopes(files)?;
270    let prediction_cache_manifest: Option<FilePredictionCacheManifest> = files
271        .get(FILE_PREDICTION_CACHE_MANIFEST_FILE)
272        .map(|bytes| parse_package_json(bytes, FILE_PREDICTION_CACHE_MANIFEST_FILE))
273        .transpose()?;
274    let artifact_manifest: Option<FileArtifactManifest> = files
275        .get(FILE_ARTIFACT_MANIFEST_FILE)
276        .map(|bytes| parse_package_json(bytes, FILE_ARTIFACT_MANIFEST_FILE))
277        .transpose()?;
278    let options = OpenLineageRunEventOptions {
279        namespace: namespace.to_string(),
280        event_time: event_time.to_string(),
281    };
282
283    build_openlineage_run_event(
284        &plan,
285        &bundle,
286        &lineage,
287        &data_envelopes,
288        prediction_cache_manifest.as_ref(),
289        artifact_manifest.as_ref(),
290        &options,
291    )
292}
293
294pub fn build_openlineage_run_event(
295    plan: &ExecutionPlan,
296    bundle: &ExecutionBundle,
297    lineage: &[LineageRecord],
298    data_envelopes: &BTreeMap<String, ExternalDataPlanEnvelope>,
299    prediction_cache_manifest: Option<&FilePredictionCacheManifest>,
300    artifact_manifest: Option<&FileArtifactManifest>,
301    options: &OpenLineageRunEventOptions,
302) -> Result<Value> {
303    validate_provenance_inputs(
304        plan,
305        bundle,
306        lineage,
307        data_envelopes,
308        prediction_cache_manifest,
309        artifact_manifest,
310    )?;
311    validate_openlineage_namespace(options.namespace.as_str())?;
312    validate_openlineage_event_time(options.event_time.as_str())?;
313
314    Ok(json!({
315        "eventType": "COMPLETE",
316        "eventTime": options.event_time.as_str(),
317        "run": {
318            "runId": openlineage_run_id(plan, bundle),
319            "facets": {
320                "dagml_reproducibility": dagml_openlineage_reproducibility_run_facet(plan, bundle),
321                "dagml_oof_safety": dagml_openlineage_oof_safety_run_facet(bundle, lineage),
322            },
323        },
324        "job": {
325            "namespace": options.namespace.as_str(),
326            "name": format!("{}::{}", plan.id, bundle.bundle_id),
327            "facets": {
328                "dagml_plan": dagml_openlineage_plan_job_facet(plan, bundle),
329            },
330        },
331        "inputs": openlineage_input_datasets(bundle, data_envelopes),
332        "outputs": openlineage_output_datasets(bundle, prediction_cache_manifest, artifact_manifest),
333        "producer": "https://github.com/GBeurier/dag-ml",
334        "schemaURL": OPENLINEAGE_RUN_EVENT_SCHEMA_URL,
335    }))
336}
337
338fn validate_provenance_inputs(
339    plan: &ExecutionPlan,
340    bundle: &ExecutionBundle,
341    lineage: &[LineageRecord],
342    data_envelopes: &BTreeMap<String, ExternalDataPlanEnvelope>,
343    prediction_cache_manifest: Option<&FilePredictionCacheManifest>,
344    artifact_manifest: Option<&FileArtifactManifest>,
345) -> Result<()> {
346    plan.validate()?;
347    bundle.validate_against_plan(plan)?;
348    if !data_envelopes.is_empty() {
349        bundle.validate_replay_envelopes(data_envelopes)?;
350    }
351    if let Some(manifest) = prediction_cache_manifest {
352        manifest.validate_against_bundle(bundle)?;
353    }
354    if let Some(manifest) = artifact_manifest {
355        manifest.validate_against_bundle(bundle)?;
356    }
357
358    let mut lineage_ids = BTreeSet::<&LineageId>::new();
359    for record in lineage {
360        record.validate()?;
361        if !plan.node_plans.contains_key(&record.node_id) {
362            return Err(DagMlError::RuntimeValidation(format!(
363                "provenance lineage `{}` references unknown node `{}`",
364                record.record_id, record.node_id
365            )));
366        }
367        if !plan
368            .controller_manifests
369            .contains_key(&record.controller_id)
370        {
371            return Err(DagMlError::RuntimeValidation(format!(
372                "provenance lineage `{}` references unknown controller `{}`",
373                record.record_id, record.controller_id
374            )));
375        }
376        if !lineage_ids.insert(&record.record_id) {
377            return Err(DagMlError::RuntimeValidation(format!(
378                "duplicate provenance lineage record `{}`",
379                record.record_id
380            )));
381        }
382    }
383    for record in lineage {
384        for input_id in &record.input_lineage {
385            if !lineage_ids.contains(input_id) {
386                return Err(DagMlError::RuntimeValidation(format!(
387                    "provenance lineage `{}` references missing input lineage `{}`",
388                    record.record_id, input_id
389                )));
390            }
391        }
392    }
393    Ok(())
394}
395
396fn build_prov_jsonld(
397    plan: &ExecutionPlan,
398    bundle: &ExecutionBundle,
399    lineage: &[LineageRecord],
400    data_envelopes: &BTreeMap<String, ExternalDataPlanEnvelope>,
401    prediction_cache_manifest: Option<&FilePredictionCacheManifest>,
402    artifact_manifest: Option<&FileArtifactManifest>,
403) -> Result<Value> {
404    let plan_entity_id = format!("dagml:execution-plan:{}", plan.id);
405    let bundle_entity_id = format!("dagml:execution-bundle:{}", bundle.bundle_id);
406    let packaging_activity_id = format!("dagml:activity:package-bundle:{}", bundle.bundle_id);
407    let coordinator_agent_id = "dagml:agent:dag-ml".to_string();
408
409    let mut entity = BTreeMap::<String, Value>::new();
410    entity.insert(
411        plan_entity_id.clone(),
412        json!({
413            "prov:type": ["prov:Entity", "dagml:ExecutionPlan"],
414            "dagml:plan_id": plan.id,
415            "dagml:graph_fingerprint": plan.graph_fingerprint,
416            "dagml:campaign_fingerprint": plan.campaign_fingerprint,
417            "dagml:controller_fingerprint": plan.controller_fingerprint,
418            "dagml:variant_count": plan.variants.len(),
419            "dagml:has_fold_set": plan.fold_set.is_some(),
420        }),
421    );
422    entity.insert(
423        bundle_entity_id.clone(),
424        json!({
425            "prov:type": ["prov:Entity", "dagml:ExecutionBundle"],
426            "dagml:bundle_id": bundle.bundle_id,
427            "dagml:schema_version": bundle.schema_version,
428            "dagml:plan_id": bundle.plan_id,
429            "dagml:selected_variant_id": bundle.selected_variant_id,
430            "dagml:graph_fingerprint": bundle.graph_fingerprint,
431            "dagml:campaign_fingerprint": bundle.campaign_fingerprint,
432            "dagml:controller_fingerprint": bundle.controller_fingerprint,
433            "dagml:unsafe_flags": bundle.unsafe_flags,
434            "dagml:selection_count": bundle.selections.len(),
435        }),
436    );
437
438    for requirement in &bundle.data_requirements {
439        let key = requirement.key();
440        entity.insert(
441            data_requirement_entity_id(&key),
442            json!({
443                "prov:type": ["prov:Entity", "dagml:DataRequirement"],
444                "dagml:requirement_key": key,
445                "dagml:node_id": requirement.node_id,
446                "dagml:input_name": requirement.input_name,
447                "dagml:schema_fingerprint": requirement.schema_fingerprint,
448                "dagml:plan_fingerprint": requirement.plan_fingerprint,
449                "dagml:relation_fingerprint": requirement.relation_fingerprint,
450                "dagml:feature_set_id": requirement.feature_set_id,
451            }),
452        );
453    }
454    for (key, envelope) in data_envelopes {
455        entity.insert(
456            data_envelope_entity_id(key),
457            json!({
458                "prov:type": ["prov:Entity", "dagml:ExternalDataPlanEnvelope"],
459                "dagml:envelope_key": key,
460                "dagml:schema_version": envelope.schema_version,
461                "dagml:schema_fingerprint": envelope.schema_fingerprint,
462                "dagml:plan_fingerprint": envelope.plan_fingerprint,
463                "dagml:relation_fingerprint": envelope.relation_fingerprint,
464            }),
465        );
466    }
467    for requirement in &bundle.prediction_requirements {
468        let key = requirement.key();
469        entity.insert(
470            prediction_requirement_entity_id(&key),
471            json!({
472                "prov:type": ["prov:Entity", "dagml:PredictionRequirement"],
473                "dagml:requirement_key": key,
474                "dagml:producer_node": requirement.producer_node,
475                "dagml:consumer_node": requirement.consumer_node,
476                "dagml:prediction_level": requirement.prediction_level,
477                "dagml:fold_ids": requirement.fold_ids,
478                "dagml:unit_ids": requirement.unit_ids,
479                "dagml:sample_ids": requirement.sample_ids,
480                "dagml:prediction_width": requirement.prediction_width,
481                "dagml:target_names": requirement.target_names,
482            }),
483        );
484    }
485    for cache in &bundle.prediction_caches {
486        entity.insert(
487            prediction_cache_entity_id(&cache.cache_id),
488            json!({
489                "prov:type": ["prov:Entity", "dagml:PredictionCache"],
490                "dagml:requirement_key": cache.requirement_key,
491                "dagml:cache_id": cache.cache_id,
492                "dagml:format": cache.format,
493                "dagml:prediction_level": cache.prediction_level,
494                "dagml:unit_ids": cache.unit_ids,
495                "dagml:block_count": cache.block_count,
496                "dagml:row_count": cache.row_count,
497                "dagml:content_fingerprint": cache.content_fingerprint,
498            }),
499        );
500    }
501    if let Some(manifest) = prediction_cache_manifest {
502        entity.insert(
503            "dagml:file:prediction-cache-manifest".to_string(),
504            json!({
505                "prov:type": ["prov:Entity", "dagml:PredictionCacheManifest"],
506                "dagml:file": FILE_PREDICTION_CACHE_MANIFEST_FILE,
507                "dagml:schema_version": manifest.schema_version,
508                "dagml:cache_count": manifest.caches.len(),
509            }),
510        );
511    }
512    for record in &bundle.refit_artifacts {
513        entity.insert(
514            artifact_entity_id(&record.artifact.id),
515            json!({
516                "prov:type": ["prov:Entity", "dagml:ModelArtifact"],
517                "dagml:artifact_id": record.artifact.id,
518                "dagml:kind": record.artifact.kind,
519                "dagml:node_id": record.node_id,
520                "dagml:controller_id": record.controller_id,
521                "dagml:backend": record.artifact.backend,
522                "dagml:uri": record.artifact.uri,
523                "dagml:content_fingerprint": record.artifact.content_fingerprint,
524                "dagml:size_bytes": record.artifact.size_bytes,
525                "dagml:plugin": record.artifact.plugin,
526                "dagml:plugin_version": record.artifact.plugin_version,
527                "dagml:params_fingerprint": record.params_fingerprint,
528                "dagml:training_loss_fingerprint": record.training_loss_fingerprint,
529                "dagml:data_requirement_keys": record.data_requirement_keys,
530                "dagml:prediction_requirement_keys": record.prediction_requirement_keys,
531            }),
532        );
533    }
534    if let Some(manifest) = artifact_manifest {
535        entity.insert(
536            "dagml:file:artifact-manifest".to_string(),
537            json!({
538                "prov:type": ["prov:Entity", "dagml:ArtifactManifest"],
539                "dagml:file": FILE_ARTIFACT_MANIFEST_FILE,
540                "dagml:schema_version": manifest.schema_version,
541                "dagml:artifact_count": manifest.artifacts.len(),
542            }),
543        );
544    }
545    for record in lineage {
546        entity.insert(
547            lineage_record_entity_id(&record.record_id),
548            json!({
549                "prov:type": ["prov:Entity", "dagml:LineageRecord"],
550                "dagml:lineage_id": record.record_id,
551                "dagml:run_id": record.run_id,
552                "dagml:node_id": record.node_id,
553                "dagml:phase": record.phase,
554                "dagml:controller_id": record.controller_id,
555                "dagml:variant_id": record.variant_id,
556                "dagml:fold_id": record.fold_id,
557                "dagml:branch_path": record.branch_path,
558                "dagml:input_lineage": record.input_lineage,
559                "dagml:artifact_refs": record
560                    .artifact_refs
561                    .iter()
562                    .map(|artifact| artifact.id.clone())
563                    .collect::<Vec<_>>(),
564            }),
565        );
566    }
567
568    let mut agent = BTreeMap::<String, Value>::new();
569    agent.insert(
570        coordinator_agent_id.clone(),
571        json!({
572            "prov:type": ["prov:Agent", "dagml:Coordinator"],
573            "dagml:name": "dag-ml",
574            "dagml:provenance_schema_version": RESEARCH_PROVENANCE_SCHEMA_VERSION,
575        }),
576    );
577    for manifest in plan.controller_manifests.values() {
578        agent.insert(
579            controller_agent_id(manifest.controller_id.as_str()),
580            json!({
581                "prov:type": ["prov:Agent", "dagml:Controller"],
582                "dagml:controller_id": manifest.controller_id,
583                "dagml:controller_version": manifest.controller_version,
584                "dagml:operator_kind": manifest.operator_kind,
585                "dagml:fit_scope": manifest.fit_scope,
586                "dagml:rng_policy": manifest.rng_policy,
587                "dagml:artifact_policy": manifest.artifact_policy,
588                "dagml:capabilities": manifest.capabilities,
589            }),
590        );
591    }
592
593    let mut activity = BTreeMap::<String, Value>::new();
594    activity.insert(
595        packaging_activity_id.clone(),
596        json!({
597            "prov:type": ["prov:Activity", "dagml:BundlePackaging"],
598            "dagml:bundle_id": bundle.bundle_id,
599            "dagml:plan_id": bundle.plan_id,
600            "dagml:selected_variant_id": bundle.selected_variant_id,
601        }),
602    );
603    for record in lineage {
604        activity.insert(
605            lineage_activity_id(record),
606            json!({
607                "prov:type": ["prov:Activity", "dagml:NodeExecution"],
608                "dagml:lineage_id": record.record_id,
609                "dagml:run_id": record.run_id,
610                "dagml:node_id": record.node_id,
611                "dagml:phase": record.phase,
612                "dagml:controller_id": record.controller_id,
613                "dagml:controller_version": record.controller_version,
614                "dagml:variant_id": record.variant_id,
615                "dagml:fold_id": record.fold_id,
616                "dagml:branch_path": record.branch_path,
617                "dagml:params_fingerprint": record.params_fingerprint,
618                "dagml:data_model_shape_fingerprint": record.data_model_shape_fingerprint,
619                "dagml:aggregation_policy_fingerprint": record.aggregation_policy_fingerprint,
620                "dagml:seed": record.seed,
621                "dagml:unsafe_flags": record.unsafe_flags,
622                "dagml:metrics": record.metrics,
623                "dagml:loss_attestations": record.loss_attestations,
624            }),
625        );
626    }
627
628    let mut used = BTreeMap::<String, Value>::new();
629    used.insert(
630        "dagml:used:bundle-plan".to_string(),
631        json!({
632            "prov:activity": packaging_activity_id,
633            "prov:entity": plan_entity_id,
634        }),
635    );
636    for record in lineage {
637        for input_id in &record.input_lineage {
638            used.insert(
639                format!("dagml:used:{}:{}", record.record_id, input_id),
640                json!({
641                    "prov:activity": lineage_activity_id(record),
642                    "prov:entity": lineage_record_entity_id(input_id),
643                    "dagml:input_lineage_id": input_id,
644                }),
645            );
646        }
647    }
648
649    let lineage_by_artifact = lineage_artifact_index(lineage);
650    let mut was_generated_by = BTreeMap::<String, Value>::new();
651    was_generated_by.insert(
652        "dagml:generated:bundle".to_string(),
653        json!({
654            "prov:entity": bundle_entity_id,
655            "prov:activity": packaging_activity_id,
656        }),
657    );
658    for record in lineage {
659        was_generated_by.insert(
660            format!("dagml:generated:lineage:{}", record.record_id),
661            json!({
662                "prov:entity": lineage_record_entity_id(&record.record_id),
663                "prov:activity": lineage_activity_id(record),
664            }),
665        );
666    }
667    for record in &bundle.refit_artifacts {
668        let activity_id = lineage_by_artifact
669            .get(&record.artifact.id)
670            .cloned()
671            .unwrap_or_else(|| packaging_activity_id.clone());
672        was_generated_by.insert(
673            format!("dagml:generated:artifact:{}", record.artifact.id),
674            json!({
675                "prov:entity": artifact_entity_id(&record.artifact.id),
676                "prov:activity": activity_id,
677            }),
678        );
679    }
680
681    let mut was_derived_from = BTreeMap::<String, Value>::new();
682    was_derived_from.insert(
683        "dagml:derived:bundle-plan".to_string(),
684        json!({
685            "prov:generatedEntity": bundle_entity_id,
686            "prov:usedEntity": plan_entity_id,
687        }),
688    );
689    for record in &bundle.refit_artifacts {
690        for key in &record.data_requirement_keys {
691            was_derived_from.insert(
692                format!("dagml:derived:{}:data:{key}", record.artifact.id),
693                json!({
694                    "prov:generatedEntity": artifact_entity_id(&record.artifact.id),
695                    "prov:usedEntity": data_requirement_entity_id(key),
696                    "dagml:refit_dependency": "data_requirement",
697                }),
698            );
699        }
700        for key in &record.prediction_requirement_keys {
701            was_derived_from.insert(
702                format!("dagml:derived:{}:prediction:{key}", record.artifact.id),
703                json!({
704                    "prov:generatedEntity": artifact_entity_id(&record.artifact.id),
705                    "prov:usedEntity": prediction_requirement_entity_id(key),
706                    "dagml:refit_dependency": "prediction_requirement",
707                    "dagml:oof_dependency": true,
708                }),
709            );
710        }
711    }
712    for cache in &bundle.prediction_caches {
713        was_derived_from.insert(
714            format!("dagml:derived:cache:{}", cache.cache_id),
715            json!({
716                "prov:generatedEntity": prediction_cache_entity_id(&cache.cache_id),
717                "prov:usedEntity": prediction_requirement_entity_id(&cache.requirement_key),
718            }),
719        );
720    }
721    for record in lineage {
722        for input_id in &record.input_lineage {
723            was_derived_from.insert(
724                format!("dagml:derived:lineage:{}:{input_id}", record.record_id),
725                json!({
726                    "prov:generatedEntity": lineage_record_entity_id(&record.record_id),
727                    "prov:usedEntity": lineage_record_entity_id(input_id),
728                    "dagml:lineage_dependency": true,
729                }),
730            );
731        }
732    }
733
734    let mut was_associated_with = BTreeMap::<String, Value>::new();
735    was_associated_with.insert(
736        "dagml:associated:bundle-packaging".to_string(),
737        json!({
738            "prov:activity": packaging_activity_id,
739            "prov:agent": coordinator_agent_id,
740        }),
741    );
742    for record in lineage {
743        was_associated_with.insert(
744            format!("dagml:associated:{}", record.record_id),
745            json!({
746                "prov:activity": lineage_activity_id(record),
747                "prov:agent": controller_agent_id(record.controller_id.as_str()),
748            }),
749        );
750    }
751
752    Ok(json!({
753        "@context": {
754            "prov": "http://www.w3.org/ns/prov#",
755            "dagml": "https://dag-ml.dev/ns#",
756        },
757        "entity": entity,
758        "activity": activity,
759        "agent": agent,
760        "used": used,
761        "wasGeneratedBy": was_generated_by,
762        "wasDerivedFrom": was_derived_from,
763        "wasAssociatedWith": was_associated_with,
764    }))
765}
766
767fn build_ro_crate_metadata(
768    plan: &ExecutionPlan,
769    bundle: &ExecutionBundle,
770    data_envelopes: &BTreeMap<String, ExternalDataPlanEnvelope>,
771    prediction_cache_manifest: Option<&FilePredictionCacheManifest>,
772    artifact_manifest: Option<&FileArtifactManifest>,
773) -> Result<Value> {
774    let mut has_part = vec![
775        json!({"@id": "execution_plan.json"}),
776        json!({"@id": "execution_bundle.json"}),
777        json!({"@id": PROV_JSONLD_FILE}),
778    ];
779    let mut graph = vec![
780        json!({
781            "@id": RO_CRATE_METADATA_FILE,
782            "@type": "CreativeWork",
783            "about": {"@id": "./"},
784            "conformsTo": {"@id": "https://w3id.org/ro/crate/1.1"},
785        }),
786        json!({
787            "@id": "./",
788            "@type": "Dataset",
789            "name": format!("DAG-ML research bundle {}", bundle.bundle_id),
790            "mainEntity": {"@id": "#workflow"},
791            "hasPart": has_part.clone(),
792            "dagml:schema_version": RESEARCH_PROVENANCE_SCHEMA_VERSION,
793            "dagml:bundle_id": bundle.bundle_id,
794            "dagml:plan_id": plan.id,
795            "dagml:unsafe_flags": bundle.unsafe_flags,
796        }),
797        json!({
798            "@id": "#workflow",
799            "@type": ["ComputationalWorkflow", "SoftwareSourceCode"],
800            "name": "DAG-ML compiled workflow",
801            "programmingLanguage": "Rust",
802            "dagml:plan_id": plan.id,
803            "dagml:graph_fingerprint": plan.graph_fingerprint,
804            "dagml:campaign_fingerprint": plan.campaign_fingerprint,
805            "dagml:controller_fingerprint": plan.controller_fingerprint,
806            "dagml:selected_variant_id": bundle.selected_variant_id,
807            "dagml:variant_count": plan.variants.len(),
808        }),
809        file_entity(
810            "execution_plan.json",
811            "DAG-ML execution plan",
812            "dagml:ExecutionPlan",
813        ),
814        file_entity(
815            "execution_bundle.json",
816            "DAG-ML execution bundle",
817            "dagml:ExecutionBundle",
818        ),
819        file_entity(PROV_JSONLD_FILE, "DAG-ML W3C PROV export", "prov:Bundle"),
820    ];
821
822    if prediction_cache_manifest.is_some() {
823        has_part.push(json!({"@id": FILE_PREDICTION_CACHE_MANIFEST_FILE}));
824        graph.push(file_entity(
825            FILE_PREDICTION_CACHE_MANIFEST_FILE,
826            "DAG-ML prediction cache manifest",
827            "dagml:PredictionCacheManifest",
828        ));
829    }
830    if artifact_manifest.is_some() {
831        has_part.push(json!({"@id": FILE_ARTIFACT_MANIFEST_FILE}));
832        graph.push(file_entity(
833            FILE_ARTIFACT_MANIFEST_FILE,
834            "DAG-ML artifact manifest",
835            "dagml:ArtifactManifest",
836        ));
837    }
838    for (key, envelope) in data_envelopes {
839        let id = format!("data_envelopes/{key}.json");
840        has_part.push(json!({"@id": id}));
841        graph.push(json!({
842            "@id": id,
843            "@type": ["File", "dagml:ExternalDataPlanEnvelope"],
844            "name": format!("DAG-ML data envelope {key}"),
845            "dagml:envelope_key": key,
846            "dagml:schema_version": envelope.schema_version,
847            "dagml:schema_fingerprint": envelope.schema_fingerprint,
848            "dagml:plan_fingerprint": envelope.plan_fingerprint,
849            "dagml:relation_fingerprint": envelope.relation_fingerprint,
850        }));
851    }
852
853    graph[1]["hasPart"] = Value::Array(has_part);
854
855    for manifest in plan.controller_manifests.values() {
856        graph.push(json!({
857            "@id": controller_agent_id(manifest.controller_id.as_str()),
858            "@type": ["SoftwareApplication", "dagml:Controller"],
859            "name": manifest.controller_id,
860            "softwareVersion": manifest.controller_version,
861            "dagml:operator_kind": manifest.operator_kind,
862            "dagml:capabilities": manifest.capabilities,
863            "dagml:artifact_policy": manifest.artifact_policy,
864        }));
865    }
866    for artifact in &bundle.refit_artifacts {
867        graph.push(json!({
868            "@id": artifact_entity_id(&artifact.artifact.id),
869            "@type": ["File", "dagml:ModelArtifact"],
870            "name": artifact.artifact.id,
871            "encodingFormat": artifact.artifact.kind,
872            "dagml:node_id": artifact.node_id,
873            "dagml:controller_id": artifact.controller_id,
874            "dagml:backend": artifact.artifact.backend,
875            "dagml:uri": artifact.artifact.uri,
876            "dagml:content_fingerprint": artifact.artifact.content_fingerprint,
877            "dagml:plugin": artifact.artifact.plugin,
878            "dagml:plugin_version": artifact.artifact.plugin_version,
879            "dagml:refit_data_requirement_keys": artifact.data_requirement_keys,
880            "dagml:refit_prediction_requirement_keys": artifact.prediction_requirement_keys,
881        }));
882    }
883
884    Ok(json!({
885        "@context": [
886            "https://w3id.org/ro/crate/1.1/context",
887            {
888                "dagml": "https://dag-ml.dev/ns#",
889                "prov": "http://www.w3.org/ns/prov#",
890            }
891        ],
892        "@graph": graph,
893    }))
894}
895
896fn file_entity(id: &str, name: &str, dagml_type: &str) -> Value {
897    json!({
898        "@id": id,
899        "@type": ["File", dagml_type],
900        "name": name,
901    })
902}
903
904fn add_json_package_file<T: Serialize + ?Sized>(
905    files: &mut BTreeMap<String, ResearchProvenancePackageFile>,
906    path: &str,
907    value: &T,
908    label: &str,
909) -> Result<()> {
910    validate_package_path(path)?;
911    let mut bytes = serde_json::to_vec_pretty(value).map_err(|err| {
912        DagMlError::RuntimeValidation(format!("failed to serialize {label}: {err}"))
913    })?;
914    bytes.push(b'\n');
915    let sha256 = sha256_hex(&bytes);
916    let previous = files.insert(
917        path.to_string(),
918        ResearchProvenancePackageFile {
919            path: path.to_string(),
920            sha256,
921            size_bytes: bytes.len(),
922            bytes,
923        },
924    );
925    if previous.is_some() {
926        return Err(DagMlError::RuntimeValidation(format!(
927            "duplicate research provenance package file `{path}`"
928        )));
929    }
930    Ok(())
931}
932
933fn validate_package_path(path: &str) -> Result<()> {
934    if path.is_empty() {
935        return Err(DagMlError::RuntimeValidation(
936            "research provenance package path is empty".to_string(),
937        ));
938    }
939    if path.starts_with('/') || path.starts_with('\\') {
940        return Err(DagMlError::RuntimeValidation(format!(
941            "research provenance package path `{path}` must be relative"
942        )));
943    }
944    if path.chars().any(char::is_control) {
945        return Err(DagMlError::RuntimeValidation(format!(
946            "research provenance package path `{path}` has control characters"
947        )));
948    }
949    for segment in path.split(['/', '\\']) {
950        if segment.is_empty() || segment == "." || segment == ".." {
951            return Err(DagMlError::RuntimeValidation(format!(
952                "research provenance package path `{path}` has an invalid path component"
953            )));
954        }
955    }
956    Ok(())
957}
958
959fn data_envelope_file_path(key: &str) -> Result<String> {
960    if key.contains(['/', '\\']) {
961        return Err(DagMlError::RuntimeValidation(format!(
962            "data envelope key `{key}` cannot be used as a research provenance package path"
963        )));
964    }
965    Ok(format!("data_envelopes/{key}.json"))
966}
967
968fn require_package_file<'a>(files: &'a BTreeMap<String, Vec<u8>>, path: &str) -> Result<&'a [u8]> {
969    files.get(path).map(Vec::as_slice).ok_or_else(|| {
970        DagMlError::RuntimeValidation(format!("research provenance package is missing `{path}`"))
971    })
972}
973
974fn parse_package_json<T: DeserializeOwned>(bytes: &[u8], path: &str) -> Result<T> {
975    serde_json::from_slice(bytes).map_err(|err| {
976        DagMlError::RuntimeValidation(format!(
977            "failed to parse research provenance package JSON `{path}`: {err}"
978        ))
979    })
980}
981
982fn parse_package_data_envelopes(
983    files: &BTreeMap<String, Vec<u8>>,
984) -> Result<BTreeMap<String, ExternalDataPlanEnvelope>> {
985    let mut envelopes = BTreeMap::new();
986    for (path, bytes) in files {
987        let Some(key) = path
988            .strip_prefix("data_envelopes/")
989            .and_then(|suffix| suffix.strip_suffix(".json"))
990        else {
991            continue;
992        };
993        if key.is_empty() || key.contains(['/', '\\']) {
994            return Err(DagMlError::RuntimeValidation(format!(
995                "research provenance data envelope path `{path}` has an invalid key"
996            )));
997        }
998        let previous = envelopes.insert(key.to_string(), parse_package_json(bytes, path)?);
999        if previous.is_some() {
1000            return Err(DagMlError::RuntimeValidation(format!(
1001                "duplicate research provenance data envelope key `{key}`"
1002            )));
1003        }
1004    }
1005    Ok(envelopes)
1006}
1007
1008fn validate_prov_jsonld_root(prov_jsonld: Value) -> Result<()> {
1009    if prov_jsonld.get("@context").is_none()
1010        || prov_jsonld.get("entity").is_none()
1011        || prov_jsonld.get("activity").is_none()
1012        || prov_jsonld.get("agent").is_none()
1013    {
1014        return Err(DagMlError::RuntimeValidation(
1015            "research provenance PROV JSON-LD root is missing required sections".to_string(),
1016        ));
1017    }
1018    Ok(())
1019}
1020
1021fn validate_openlineage_namespace(namespace: &str) -> Result<()> {
1022    if namespace.trim().is_empty() {
1023        return Err(DagMlError::RuntimeValidation(
1024            "OpenLineage namespace must not be empty".to_string(),
1025        ));
1026    }
1027    if namespace.chars().any(char::is_control) {
1028        return Err(DagMlError::RuntimeValidation(
1029            "OpenLineage namespace contains control characters".to_string(),
1030        ));
1031    }
1032    Ok(())
1033}
1034
1035fn validate_openlineage_event_time(event_time: &str) -> Result<()> {
1036    if event_time.trim().is_empty() || !event_time.contains('T') {
1037        return Err(DagMlError::RuntimeValidation(
1038            "OpenLineage event_time must be a non-empty RFC3339-like timestamp".to_string(),
1039        ));
1040    }
1041    if event_time.chars().any(char::is_control) {
1042        return Err(DagMlError::RuntimeValidation(
1043            "OpenLineage event_time contains control characters".to_string(),
1044        ));
1045    }
1046    Ok(())
1047}
1048
1049fn dagml_openlineage_reproducibility_run_facet(
1050    plan: &ExecutionPlan,
1051    bundle: &ExecutionBundle,
1052) -> Value {
1053    json!({
1054        "_schemaURL": format!("{DAGML_OPENLINEAGE_FACET_SCHEMA_URL}#/$defs/DagmlReproducibilityRunFacet"),
1055        "plan_id": plan.id,
1056        "bundle_id": bundle.bundle_id,
1057        "graph_fingerprint": bundle.graph_fingerprint,
1058        "campaign_fingerprint": bundle.campaign_fingerprint,
1059        "controller_fingerprint": bundle.controller_fingerprint,
1060        "selected_variant_id": bundle.selected_variant_id,
1061        "variant_count": plan.variants.len(),
1062        "unsafe_flags": bundle.unsafe_flags,
1063    })
1064}
1065
1066fn dagml_openlineage_oof_safety_run_facet(
1067    bundle: &ExecutionBundle,
1068    lineage: &[LineageRecord],
1069) -> Value {
1070    json!({
1071        "_schemaURL": format!("{DAGML_OPENLINEAGE_FACET_SCHEMA_URL}#/$defs/DagmlOofSafetyRunFacet"),
1072        "prediction_requirement_count": bundle.prediction_requirements.len(),
1073        "prediction_cache_count": bundle.prediction_caches.len(),
1074        "lineage_record_count": lineage.len(),
1075        "requires_oof_prediction_count": bundle.prediction_requirements.len(),
1076        "refit_artifact_count": bundle.refit_artifacts.len(),
1077    })
1078}
1079
1080fn dagml_openlineage_plan_job_facet(plan: &ExecutionPlan, bundle: &ExecutionBundle) -> Value {
1081    json!({
1082        "_schemaURL": format!("{DAGML_OPENLINEAGE_FACET_SCHEMA_URL}#/$defs/DagmlPlanJobFacet"),
1083        "plan_id": plan.id,
1084        "bundle_id": bundle.bundle_id,
1085        "node_count": plan.node_plans.len(),
1086        "controller_count": plan.controller_manifests.len(),
1087        "has_fold_set": plan.fold_set.is_some(),
1088        "selected_variant_id": bundle.selected_variant_id,
1089    })
1090}
1091
1092fn openlineage_input_datasets(
1093    bundle: &ExecutionBundle,
1094    data_envelopes: &BTreeMap<String, ExternalDataPlanEnvelope>,
1095) -> Vec<Value> {
1096    bundle
1097        .data_requirements
1098        .iter()
1099        .map(|requirement| {
1100            let key = requirement.key();
1101            let envelope = data_envelopes.get(&key);
1102            json!({
1103                "namespace": "dagml:data-requirement",
1104                "name": key,
1105                "facets": {
1106                    "dagml_contract": {
1107                        "_schemaURL": format!("{DAGML_OPENLINEAGE_FACET_SCHEMA_URL}#/$defs/DagmlDatasetContractFacet"),
1108                        "node_id": requirement.node_id,
1109                        "input_name": requirement.input_name,
1110                        "schema_fingerprint": requirement.schema_fingerprint,
1111                        "plan_fingerprint": requirement.plan_fingerprint,
1112                        "relation_fingerprint": requirement.relation_fingerprint,
1113                        "feature_set_id": requirement.feature_set_id,
1114                        "envelope_schema_fingerprint": envelope.map(|envelope| envelope.schema_fingerprint.clone()),
1115                        "envelope_plan_fingerprint": envelope.map(|envelope| envelope.plan_fingerprint.clone()),
1116                    }
1117                }
1118            })
1119        })
1120        .collect()
1121}
1122
1123fn openlineage_output_datasets(
1124    bundle: &ExecutionBundle,
1125    prediction_cache_manifest: Option<&FilePredictionCacheManifest>,
1126    artifact_manifest: Option<&FileArtifactManifest>,
1127) -> Vec<Value> {
1128    let mut outputs = vec![json!({
1129        "namespace": "dagml:bundle",
1130        "name": bundle.bundle_id,
1131        "facets": {
1132            "dagml_contract": {
1133                "_schemaURL": format!("{DAGML_OPENLINEAGE_FACET_SCHEMA_URL}#/$defs/DagmlDatasetContractFacet"),
1134                "schema_version": bundle.schema_version,
1135                "plan_id": bundle.plan_id,
1136                "selected_variant_id": bundle.selected_variant_id,
1137                "graph_fingerprint": bundle.graph_fingerprint,
1138                "campaign_fingerprint": bundle.campaign_fingerprint,
1139                "controller_fingerprint": bundle.controller_fingerprint,
1140            }
1141        }
1142    })];
1143    for cache in &bundle.prediction_caches {
1144        outputs.push(json!({
1145            "namespace": "dagml:prediction-cache",
1146            "name": cache.cache_id,
1147            "facets": {
1148                "dagml_contract": {
1149                    "_schemaURL": format!("{DAGML_OPENLINEAGE_FACET_SCHEMA_URL}#/$defs/DagmlDatasetContractFacet"),
1150                    "requirement_key": cache.requirement_key,
1151                    "prediction_level": cache.prediction_level,
1152                    "row_count": cache.row_count,
1153                    "block_count": cache.block_count,
1154                    "content_fingerprint": cache.content_fingerprint,
1155                    "has_file_manifest": prediction_cache_manifest.is_some(),
1156                }
1157            }
1158        }));
1159    }
1160    for artifact in &bundle.refit_artifacts {
1161        outputs.push(json!({
1162            "namespace": "dagml:artifact",
1163            "name": artifact.artifact.id,
1164            "facets": {
1165                "dagml_contract": {
1166                    "_schemaURL": format!("{DAGML_OPENLINEAGE_FACET_SCHEMA_URL}#/$defs/DagmlDatasetContractFacet"),
1167                    "node_id": artifact.node_id,
1168                    "controller_id": artifact.controller_id,
1169                    "backend": artifact.artifact.backend,
1170                    "uri": artifact.artifact.uri,
1171                    "content_fingerprint": artifact.artifact.content_fingerprint,
1172                    "plugin": artifact.artifact.plugin,
1173                    "plugin_version": artifact.artifact.plugin_version,
1174                    "data_requirement_keys": artifact.data_requirement_keys,
1175                    "prediction_requirement_keys": artifact.prediction_requirement_keys,
1176                    "has_file_manifest": artifact_manifest.is_some(),
1177                }
1178            }
1179        }));
1180    }
1181    outputs
1182}
1183
1184fn openlineage_run_id(plan: &ExecutionPlan, bundle: &ExecutionBundle) -> String {
1185    let input = format!("dag-ml/openlineage/run/{}/{}", plan.id, bundle.bundle_id);
1186    let digest = Sha256::digest(input.as_bytes());
1187    let mut bytes = [0u8; 16];
1188    bytes.copy_from_slice(&digest[..16]);
1189    bytes[6] = (bytes[6] & 0x0f) | 0x50;
1190    bytes[8] = (bytes[8] & 0x3f) | 0x80;
1191    format!(
1192        "{:02x}{:02x}{:02x}{:02x}-{:02x}{:02x}-{:02x}{:02x}-{:02x}{:02x}-{:02x}{:02x}{:02x}{:02x}{:02x}{:02x}",
1193        bytes[0],
1194        bytes[1],
1195        bytes[2],
1196        bytes[3],
1197        bytes[4],
1198        bytes[5],
1199        bytes[6],
1200        bytes[7],
1201        bytes[8],
1202        bytes[9],
1203        bytes[10],
1204        bytes[11],
1205        bytes[12],
1206        bytes[13],
1207        bytes[14],
1208        bytes[15],
1209    )
1210}
1211
1212fn validate_ro_crate_package_checksums(
1213    ro_crate_metadata: &Value,
1214    files: &BTreeMap<String, Vec<u8>>,
1215) -> Result<usize> {
1216    let graph = ro_crate_metadata
1217        .get("@graph")
1218        .and_then(Value::as_array)
1219        .ok_or_else(|| {
1220            DagMlError::RuntimeValidation("RO-Crate metadata has no @graph array".to_string())
1221        })?;
1222    let root = graph
1223        .iter()
1224        .find(|entry| entry.get("@id").and_then(Value::as_str) == Some("./"))
1225        .ok_or_else(|| {
1226            DagMlError::RuntimeValidation("RO-Crate metadata has no root dataset".to_string())
1227        })?;
1228    if root.get("dagml:schema_version").and_then(Value::as_u64)
1229        != Some(RESEARCH_PROVENANCE_SCHEMA_VERSION as u64)
1230    {
1231        return Err(DagMlError::RuntimeValidation(format!(
1232            "RO-Crate root has unsupported dagml:schema_version, expected {}",
1233            RESEARCH_PROVENANCE_SCHEMA_VERSION
1234        )));
1235    }
1236    let root_has_part = root
1237        .get("hasPart")
1238        .and_then(Value::as_array)
1239        .ok_or_else(|| {
1240            DagMlError::RuntimeValidation("RO-Crate root hasPart is not an array".to_string())
1241        })?;
1242    let root_has_part_ids = root_has_part
1243        .iter()
1244        .filter_map(|entry| entry.get("@id").and_then(Value::as_str))
1245        .collect::<BTreeSet<_>>();
1246
1247    let mut checksummed = 0;
1248    for (path, bytes) in files {
1249        if path == RO_CRATE_METADATA_FILE {
1250            continue;
1251        }
1252        if !root_has_part_ids.contains(path.as_str()) {
1253            return Err(DagMlError::RuntimeValidation(format!(
1254                "RO-Crate root does not list package file `{path}` in hasPart"
1255            )));
1256        }
1257        let entry = graph
1258            .iter()
1259            .find(|entry| entry.get("@id").and_then(Value::as_str) == Some(path.as_str()))
1260            .ok_or_else(|| {
1261                DagMlError::RuntimeValidation(format!(
1262                    "RO-Crate metadata is missing package file entry `{path}`"
1263                ))
1264            })?;
1265        let expected_sha256 = sha256_hex(bytes);
1266        let declared_sha256 = entry.get("sha256").and_then(Value::as_str).ok_or_else(|| {
1267            DagMlError::RuntimeValidation(format!(
1268                "RO-Crate package file `{path}` is missing sha256"
1269            ))
1270        })?;
1271        if declared_sha256 != expected_sha256 {
1272            return Err(DagMlError::RuntimeValidation(format!(
1273                "RO-Crate package file `{path}` sha256 mismatch"
1274            )));
1275        }
1276        let declared_dagml_sha256 = entry
1277            .get("dagml:sha256")
1278            .and_then(Value::as_str)
1279            .ok_or_else(|| {
1280                DagMlError::RuntimeValidation(format!(
1281                    "RO-Crate package file `{path}` is missing dagml:sha256"
1282                ))
1283            })?;
1284        if declared_dagml_sha256 != declared_sha256 {
1285            return Err(DagMlError::RuntimeValidation(format!(
1286                "RO-Crate package file `{path}` has inconsistent checksum fields"
1287            )));
1288        }
1289        if entry.get("contentSize").and_then(Value::as_u64) != Some(bytes.len() as u64) {
1290            return Err(DagMlError::RuntimeValidation(format!(
1291                "RO-Crate package file `{path}` contentSize mismatch"
1292            )));
1293        }
1294        if entry.get("encodingFormat").and_then(Value::as_str) != Some("application/json") {
1295            return Err(DagMlError::RuntimeValidation(format!(
1296                "RO-Crate package file `{path}` must declare application/json"
1297            )));
1298        }
1299        checksummed += 1;
1300    }
1301
1302    for has_part_id in root_has_part_ids {
1303        if has_part_id != RO_CRATE_METADATA_FILE && !files.contains_key(has_part_id) {
1304            return Err(DagMlError::RuntimeValidation(format!(
1305                "RO-Crate root references missing package file `{has_part_id}`"
1306            )));
1307        }
1308    }
1309
1310    Ok(checksummed)
1311}
1312
1313fn annotate_ro_crate_package_files(
1314    ro_crate_metadata: &mut Value,
1315    files: &BTreeMap<String, ResearchProvenancePackageFile>,
1316) -> Result<()> {
1317    let graph = ro_crate_metadata
1318        .get_mut("@graph")
1319        .and_then(Value::as_array_mut)
1320        .ok_or_else(|| {
1321            DagMlError::RuntimeValidation("RO-Crate metadata has no @graph array".to_string())
1322        })?;
1323
1324    let mut existing_ids = graph
1325        .iter()
1326        .filter_map(|entry| entry.get("@id").and_then(Value::as_str).map(str::to_string))
1327        .collect::<BTreeSet<_>>();
1328    for file in files.values() {
1329        if !existing_ids.contains(&file.path) {
1330            graph.push(file_entity(
1331                &file.path,
1332                &format!("DAG-ML contract file {}", file.path),
1333                "dagml:ContractArtifact",
1334            ));
1335            existing_ids.insert(file.path.clone());
1336        }
1337    }
1338
1339    for entry in graph.iter_mut() {
1340        let Some(id) = entry.get("@id").and_then(Value::as_str).map(str::to_string) else {
1341            continue;
1342        };
1343        let Some(file) = files.get(id.as_str()) else {
1344            continue;
1345        };
1346        let object = entry.as_object_mut().ok_or_else(|| {
1347            DagMlError::RuntimeValidation(format!("RO-Crate graph entry `{id}` is not an object"))
1348        })?;
1349        object.insert("encodingFormat".to_string(), json!("application/json"));
1350        object.insert("contentSize".to_string(), json!(file.size_bytes));
1351        object.insert("sha256".to_string(), json!(file.sha256));
1352        object.insert("dagml:sha256".to_string(), json!(file.sha256));
1353    }
1354
1355    let root = graph
1356        .iter_mut()
1357        .find(|entry| entry.get("@id") == Some(&json!("./")))
1358        .ok_or_else(|| {
1359            DagMlError::RuntimeValidation("RO-Crate metadata has no root dataset".to_string())
1360        })?;
1361    let root_object = root.as_object_mut().ok_or_else(|| {
1362        DagMlError::RuntimeValidation("RO-Crate root dataset is not an object".to_string())
1363    })?;
1364    let has_part = root_object
1365        .entry("hasPart".to_string())
1366        .or_insert_with(|| Value::Array(Vec::new()));
1367    let has_part = has_part.as_array_mut().ok_or_else(|| {
1368        DagMlError::RuntimeValidation("RO-Crate root hasPart is not an array".to_string())
1369    })?;
1370    let mut has_part_ids = has_part
1371        .iter()
1372        .filter_map(|entry| entry.get("@id").and_then(Value::as_str).map(str::to_string))
1373        .collect::<BTreeSet<_>>();
1374    for path in files.keys() {
1375        if has_part_ids.insert(path.clone()) {
1376            has_part.push(json!({"@id": path}));
1377        }
1378    }
1379    Ok(())
1380}
1381
1382fn sha256_hex(bytes: &[u8]) -> String {
1383    let digest = Sha256::digest(bytes);
1384    let mut out = String::with_capacity(digest.len() * 2);
1385    for byte in digest {
1386        use std::fmt::Write;
1387        write!(&mut out, "{byte:02x}").expect("writing to string cannot fail");
1388    }
1389    out
1390}
1391
1392fn lineage_artifact_index(lineage: &[LineageRecord]) -> BTreeMap<ArtifactId, String> {
1393    let mut index = BTreeMap::new();
1394    for record in lineage {
1395        for artifact in &record.artifact_refs {
1396            index.insert(artifact.id.clone(), lineage_activity_id(record));
1397        }
1398    }
1399    index
1400}
1401
1402fn lineage_activity_id(record: &LineageRecord) -> String {
1403    format!("dagml:activity:{}", record.record_id)
1404}
1405
1406fn controller_agent_id(controller_id: &str) -> String {
1407    format!("dagml:controller:{controller_id}")
1408}
1409
1410fn artifact_entity_id(artifact_id: &ArtifactId) -> String {
1411    format!("dagml:artifact:{artifact_id}")
1412}
1413
1414fn lineage_record_entity_id(lineage_id: &LineageId) -> String {
1415    format!("dagml:lineage-record:{lineage_id}")
1416}
1417
1418fn data_requirement_entity_id(key: &str) -> String {
1419    format!("dagml:data-requirement:{key}")
1420}
1421
1422fn data_envelope_entity_id(key: &str) -> String {
1423    format!("dagml:data-envelope:{key}")
1424}
1425
1426fn prediction_requirement_entity_id(key: &str) -> String {
1427    format!("dagml:prediction-requirement:{key}")
1428}
1429
1430fn prediction_cache_entity_id(cache_id: &str) -> String {
1431    format!("dagml:prediction-cache:{cache_id}")
1432}
1433
1434#[cfg(test)]
1435mod tests {
1436    use super::*;
1437    use crate::controller::{ControllerManifest, ControllerRegistry};
1438    use crate::ids::{ControllerId, LineageId, NodeId, RunId};
1439    use crate::plan::build_execution_plan;
1440    use crate::{CampaignSpec, GraphSpec, Phase};
1441
1442    fn fixture_plan() -> ExecutionPlan {
1443        let graph: GraphSpec =
1444            serde_json::from_str(include_str!("../tests/fixtures/package/minimal_graph.json"))
1445                .unwrap();
1446        let campaign: CampaignSpec = serde_json::from_str(include_str!(
1447            "../tests/fixtures/package/campaign_oof_generation.json"
1448        ))
1449        .unwrap();
1450        let manifests: Vec<ControllerManifest> = serde_json::from_str(include_str!(
1451            "../tests/fixtures/package/controller_manifests.json"
1452        ))
1453        .unwrap();
1454        let mut registry = ControllerRegistry::new();
1455        for manifest in manifests {
1456            registry.register(manifest).unwrap();
1457        }
1458        build_execution_plan("plan:cli.bundle", graph, campaign, &registry).unwrap()
1459    }
1460
1461    fn fixture_bundle() -> ExecutionBundle {
1462        serde_json::from_str(include_str!(
1463            "../tests/fixtures/package/provenance/execution_bundle_minimal.json"
1464        ))
1465        .unwrap()
1466    }
1467
1468    fn fixture_lineage(plan: &ExecutionPlan) -> LineageRecord {
1469        let node_id = NodeId::new("model:base").unwrap();
1470        let node_plan = plan.node_plans.get(&node_id).unwrap();
1471        LineageRecord {
1472            record_id: LineageId::new("lineage:test:model:base").unwrap(),
1473            run_id: RunId::new("run:provenance").unwrap(),
1474            node_id,
1475            phase: Phase::Refit,
1476            controller_id: node_plan.controller_id.clone(),
1477            controller_version: node_plan.controller_version.clone(),
1478            variant_id: plan
1479                .variants
1480                .first()
1481                .map(|variant| variant.variant_id.clone()),
1482            fold_id: None,
1483            branch_path: Vec::new(),
1484            input_lineage: Vec::new(),
1485            artifact_refs: vec![fixture_bundle().refit_artifacts[0].artifact.clone()],
1486            params_fingerprint: node_plan.params_fingerprint.clone(),
1487            data_model_shape_fingerprint: None,
1488            aggregation_policy_fingerprint: None,
1489            seed: Some(42),
1490            unsafe_flags: BTreeSet::new(),
1491            metrics: BTreeMap::new(),
1492            loss_attestations: Vec::new(),
1493            early_stopping_records: Vec::new(),
1494        }
1495    }
1496
1497    #[test]
1498    fn research_provenance_export_contains_prov_and_ro_crate_contracts() {
1499        let plan = fixture_plan();
1500        let bundle = fixture_bundle();
1501        let lineage = vec![fixture_lineage(&plan)];
1502        let export = build_research_provenance_export(
1503            &plan,
1504            &bundle,
1505            &lineage,
1506            &BTreeMap::new(),
1507            None,
1508            None,
1509        )
1510        .unwrap();
1511
1512        assert_eq!(export.schema_version, RESEARCH_PROVENANCE_SCHEMA_VERSION);
1513        assert!(export.prov_jsonld["@context"]["prov"]
1514            .as_str()
1515            .unwrap()
1516            .contains("prov"));
1517        assert!(export.prov_jsonld["activity"]
1518            .as_object()
1519            .unwrap()
1520            .contains_key("dagml:activity:lineage:test:model:base"));
1521        assert!(export.prov_jsonld["agent"]
1522            .as_object()
1523            .unwrap()
1524            .contains_key("dagml:controller:controller:model.mock"));
1525        assert!(export.prov_jsonld["entity"]
1526            .as_object()
1527            .unwrap()
1528            .contains_key("dagml:artifact:artifact:model:base:refit"));
1529
1530        let graph = export.ro_crate_metadata["@graph"].as_array().unwrap();
1531        assert!(graph
1532            .iter()
1533            .any(|entry| entry["@type"].to_string().contains("ComputationalWorkflow")));
1534        assert!(graph
1535            .iter()
1536            .any(|entry| entry["@id"] == json!("lineage.prov.jsonld")));
1537        assert!(graph
1538            .iter()
1539            .any(|entry| entry["@id"] == json!("execution_bundle.json")));
1540    }
1541
1542    #[test]
1543    fn research_provenance_package_contains_contract_files_and_checksums() {
1544        let plan = fixture_plan();
1545        let bundle = fixture_bundle();
1546        let lineage = vec![fixture_lineage(&plan)];
1547        let package = build_research_provenance_package(
1548            &plan,
1549            &bundle,
1550            &lineage,
1551            &BTreeMap::new(),
1552            None,
1553            None,
1554        )
1555        .unwrap();
1556
1557        for path in [
1558            EXECUTION_PLAN_FILE,
1559            EXECUTION_BUNDLE_FILE,
1560            LINEAGE_RECORDS_FILE,
1561            PROV_JSONLD_FILE,
1562            RO_CRATE_METADATA_FILE,
1563        ] {
1564            assert!(
1565                package.files.contains_key(path),
1566                "package is missing {path}"
1567            );
1568        }
1569        for (path, file) in &package.files {
1570            assert_eq!(file.path, *path);
1571            assert_eq!(file.sha256.len(), 64, "invalid sha256 for {path}");
1572            assert!(file.size_bytes > 0, "empty package file {path}");
1573            assert_eq!(file.size_bytes, file.bytes.len());
1574        }
1575
1576        let ro_crate_file = package.files.get(RO_CRATE_METADATA_FILE).unwrap();
1577        let ro_crate_metadata: Value = serde_json::from_slice(&ro_crate_file.bytes).unwrap();
1578        let graph = ro_crate_metadata["@graph"].as_array().unwrap();
1579        for path in [
1580            EXECUTION_PLAN_FILE,
1581            EXECUTION_BUNDLE_FILE,
1582            LINEAGE_RECORDS_FILE,
1583            PROV_JSONLD_FILE,
1584        ] {
1585            let entry = graph
1586                .iter()
1587                .find(|entry| entry["@id"] == json!(path))
1588                .unwrap_or_else(|| panic!("RO-Crate metadata is missing file entry {path}"));
1589            assert_eq!(entry["sha256"].as_str().map(str::len), Some(64));
1590            assert_eq!(entry["dagml:sha256"].as_str(), entry["sha256"].as_str());
1591            assert_eq!(entry["encodingFormat"], json!("application/json"));
1592            assert!(entry["contentSize"].as_u64().unwrap() > 0);
1593        }
1594        let root = graph
1595            .iter()
1596            .find(|entry| entry["@id"] == json!("./"))
1597            .expect("RO-Crate root dataset is present");
1598        let has_part = root["hasPart"].as_array().unwrap();
1599        assert!(has_part
1600            .iter()
1601            .any(|entry| entry["@id"] == json!(LINEAGE_RECORDS_FILE)));
1602    }
1603
1604    #[test]
1605    fn research_provenance_package_validation_reopens_exported_contracts() {
1606        let plan = fixture_plan();
1607        let bundle = fixture_bundle();
1608        let lineage = vec![fixture_lineage(&plan)];
1609        let package = build_research_provenance_package(
1610            &plan,
1611            &bundle,
1612            &lineage,
1613            &BTreeMap::new(),
1614            None,
1615            None,
1616        )
1617        .unwrap();
1618        let files = package
1619            .files
1620            .iter()
1621            .map(|(path, file)| (path.clone(), file.bytes.clone()))
1622            .collect::<BTreeMap<_, _>>();
1623
1624        let validation = validate_research_provenance_package_files(&files).unwrap();
1625
1626        assert_eq!(
1627            validation.schema_version,
1628            RESEARCH_PROVENANCE_SCHEMA_VERSION
1629        );
1630        assert_eq!(validation.plan_id, plan.id.to_string());
1631        assert_eq!(validation.bundle_id, bundle.bundle_id.to_string());
1632        assert_eq!(validation.file_count, package.files.len());
1633        assert_eq!(validation.checksummed_file_count, package.files.len() - 1);
1634        assert_eq!(validation.lineage_record_count, 1);
1635    }
1636
1637    #[test]
1638    fn research_provenance_package_validation_refuses_tampered_file() {
1639        let plan = fixture_plan();
1640        let bundle = fixture_bundle();
1641        let package =
1642            build_research_provenance_package(&plan, &bundle, &[], &BTreeMap::new(), None, None)
1643                .unwrap();
1644        let mut files = package
1645            .files
1646            .iter()
1647            .map(|(path, file)| (path.clone(), file.bytes.clone()))
1648            .collect::<BTreeMap<_, _>>();
1649        files.insert(LINEAGE_RECORDS_FILE.to_string(), b"[]\n\n".to_vec());
1650
1651        let error = validate_research_provenance_package_files(&files)
1652            .unwrap_err()
1653            .to_string();
1654
1655        assert!(
1656            error.contains("sha256 mismatch"),
1657            "unexpected error: {error}"
1658        );
1659    }
1660
1661    #[test]
1662    fn openlineage_export_uses_validated_research_package() {
1663        let plan = fixture_plan();
1664        let bundle = fixture_bundle();
1665        let lineage = vec![fixture_lineage(&plan)];
1666        let package = build_research_provenance_package(
1667            &plan,
1668            &bundle,
1669            &lineage,
1670            &BTreeMap::new(),
1671            None,
1672            None,
1673        )
1674        .unwrap();
1675        let files = package
1676            .files
1677            .iter()
1678            .map(|(path, file)| (path.clone(), file.bytes.clone()))
1679            .collect::<BTreeMap<_, _>>();
1680
1681        let event = build_openlineage_run_event_from_package_files(
1682            &files,
1683            "dag-ml-test",
1684            "2026-05-27T00:00:00Z",
1685        )
1686        .unwrap();
1687
1688        assert_eq!(event["eventType"], json!("COMPLETE"));
1689        assert_eq!(event["schemaURL"], json!(OPENLINEAGE_RUN_EVENT_SCHEMA_URL));
1690        assert_eq!(event["job"]["namespace"], json!("dag-ml-test"));
1691        assert_eq!(event["run"]["runId"].as_str().map(str::len), Some(36));
1692        assert!(
1693            event["run"]["facets"]["dagml_reproducibility"]["graph_fingerprint"]
1694                .as_str()
1695                .is_some()
1696        );
1697        assert!(!event["inputs"].as_array().unwrap().is_empty());
1698        assert!(event["outputs"]
1699            .as_array()
1700            .unwrap()
1701            .iter()
1702            .any(|output| output["namespace"] == json!("dagml:bundle")));
1703    }
1704
1705    #[test]
1706    fn research_provenance_export_refuses_unknown_lineage_node() {
1707        let plan = fixture_plan();
1708        let bundle = fixture_bundle();
1709        let mut lineage = fixture_lineage(&plan);
1710        lineage.node_id = NodeId::new("model:missing").unwrap();
1711
1712        let error = build_research_provenance_export(
1713            &plan,
1714            &bundle,
1715            &[lineage],
1716            &BTreeMap::new(),
1717            None,
1718            None,
1719        )
1720        .unwrap_err()
1721        .to_string();
1722
1723        assert!(error.contains("unknown node"), "unexpected error: {error}");
1724    }
1725
1726    #[test]
1727    fn research_provenance_export_refuses_mismatched_artifact_manifest() {
1728        let plan = fixture_plan();
1729        let bundle = fixture_bundle();
1730        let mut manifest = FileArtifactManifest {
1731            bundle_id: bundle.bundle_id.clone(),
1732            schema_version: crate::runtime::FILE_ARTIFACT_MANIFEST_SCHEMA_VERSION,
1733            artifacts: Vec::new(),
1734        };
1735        manifest.bundle_id = crate::ids::BundleId::new("bundle:wrong").unwrap();
1736
1737        let error = build_research_provenance_export(
1738            &plan,
1739            &bundle,
1740            &[],
1741            &BTreeMap::new(),
1742            None,
1743            Some(&manifest),
1744        )
1745        .unwrap_err()
1746        .to_string();
1747
1748        assert!(
1749            error.contains("does not match bundle"),
1750            "unexpected error: {error}"
1751        );
1752    }
1753
1754    #[test]
1755    fn research_provenance_export_refuses_unknown_lineage_controller() {
1756        let plan = fixture_plan();
1757        let bundle = fixture_bundle();
1758        let mut lineage = fixture_lineage(&plan);
1759        lineage.controller_id = ControllerId::new("controller:missing").unwrap();
1760
1761        let error = build_research_provenance_export(
1762            &plan,
1763            &bundle,
1764            &[lineage],
1765            &BTreeMap::new(),
1766            None,
1767            None,
1768        )
1769        .unwrap_err()
1770        .to_string();
1771
1772        assert!(
1773            error.contains("unknown controller"),
1774            "unexpected error: {error}"
1775        );
1776    }
1777}