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!("../../../examples/minimal_graph.json")).unwrap();
1445        let campaign: CampaignSpec = serde_json::from_str(include_str!(
1446            "../../../examples/campaign_oof_generation.json"
1447        ))
1448        .unwrap();
1449        let manifests: Vec<ControllerManifest> =
1450            serde_json::from_str(include_str!("../../../examples/controller_manifests.json"))
1451                .unwrap();
1452        let mut registry = ControllerRegistry::new();
1453        for manifest in manifests {
1454            registry.register(manifest).unwrap();
1455        }
1456        build_execution_plan("plan:cli.bundle", graph, campaign, &registry).unwrap()
1457    }
1458
1459    fn fixture_bundle() -> ExecutionBundle {
1460        serde_json::from_str(include_str!(
1461            "../../../examples/generated/execution_bundle_minimal.json"
1462        ))
1463        .unwrap()
1464    }
1465
1466    fn fixture_lineage(plan: &ExecutionPlan) -> LineageRecord {
1467        let node_id = NodeId::new("model:base").unwrap();
1468        let node_plan = plan.node_plans.get(&node_id).unwrap();
1469        LineageRecord {
1470            record_id: LineageId::new("lineage:test:model:base").unwrap(),
1471            run_id: RunId::new("run:provenance").unwrap(),
1472            node_id,
1473            phase: Phase::Refit,
1474            controller_id: node_plan.controller_id.clone(),
1475            controller_version: node_plan.controller_version.clone(),
1476            variant_id: plan
1477                .variants
1478                .first()
1479                .map(|variant| variant.variant_id.clone()),
1480            fold_id: None,
1481            branch_path: Vec::new(),
1482            input_lineage: Vec::new(),
1483            artifact_refs: vec![fixture_bundle().refit_artifacts[0].artifact.clone()],
1484            params_fingerprint: node_plan.params_fingerprint.clone(),
1485            data_model_shape_fingerprint: None,
1486            aggregation_policy_fingerprint: None,
1487            seed: Some(42),
1488            unsafe_flags: BTreeSet::new(),
1489            metrics: BTreeMap::new(),
1490            loss_attestations: Vec::new(),
1491            early_stopping_records: Vec::new(),
1492        }
1493    }
1494
1495    #[test]
1496    fn research_provenance_export_contains_prov_and_ro_crate_contracts() {
1497        let plan = fixture_plan();
1498        let bundle = fixture_bundle();
1499        let lineage = vec![fixture_lineage(&plan)];
1500        let export = build_research_provenance_export(
1501            &plan,
1502            &bundle,
1503            &lineage,
1504            &BTreeMap::new(),
1505            None,
1506            None,
1507        )
1508        .unwrap();
1509
1510        assert_eq!(export.schema_version, RESEARCH_PROVENANCE_SCHEMA_VERSION);
1511        assert!(export.prov_jsonld["@context"]["prov"]
1512            .as_str()
1513            .unwrap()
1514            .contains("prov"));
1515        assert!(export.prov_jsonld["activity"]
1516            .as_object()
1517            .unwrap()
1518            .contains_key("dagml:activity:lineage:test:model:base"));
1519        assert!(export.prov_jsonld["agent"]
1520            .as_object()
1521            .unwrap()
1522            .contains_key("dagml:controller:controller:model.mock"));
1523        assert!(export.prov_jsonld["entity"]
1524            .as_object()
1525            .unwrap()
1526            .contains_key("dagml:artifact:artifact:model:base:refit"));
1527
1528        let graph = export.ro_crate_metadata["@graph"].as_array().unwrap();
1529        assert!(graph
1530            .iter()
1531            .any(|entry| entry["@type"].to_string().contains("ComputationalWorkflow")));
1532        assert!(graph
1533            .iter()
1534            .any(|entry| entry["@id"] == json!("lineage.prov.jsonld")));
1535        assert!(graph
1536            .iter()
1537            .any(|entry| entry["@id"] == json!("execution_bundle.json")));
1538    }
1539
1540    #[test]
1541    fn research_provenance_package_contains_contract_files_and_checksums() {
1542        let plan = fixture_plan();
1543        let bundle = fixture_bundle();
1544        let lineage = vec![fixture_lineage(&plan)];
1545        let package = build_research_provenance_package(
1546            &plan,
1547            &bundle,
1548            &lineage,
1549            &BTreeMap::new(),
1550            None,
1551            None,
1552        )
1553        .unwrap();
1554
1555        for path in [
1556            EXECUTION_PLAN_FILE,
1557            EXECUTION_BUNDLE_FILE,
1558            LINEAGE_RECORDS_FILE,
1559            PROV_JSONLD_FILE,
1560            RO_CRATE_METADATA_FILE,
1561        ] {
1562            assert!(
1563                package.files.contains_key(path),
1564                "package is missing {path}"
1565            );
1566        }
1567        for (path, file) in &package.files {
1568            assert_eq!(file.path, *path);
1569            assert_eq!(file.sha256.len(), 64, "invalid sha256 for {path}");
1570            assert!(file.size_bytes > 0, "empty package file {path}");
1571            assert_eq!(file.size_bytes, file.bytes.len());
1572        }
1573
1574        let ro_crate_file = package.files.get(RO_CRATE_METADATA_FILE).unwrap();
1575        let ro_crate_metadata: Value = serde_json::from_slice(&ro_crate_file.bytes).unwrap();
1576        let graph = ro_crate_metadata["@graph"].as_array().unwrap();
1577        for path in [
1578            EXECUTION_PLAN_FILE,
1579            EXECUTION_BUNDLE_FILE,
1580            LINEAGE_RECORDS_FILE,
1581            PROV_JSONLD_FILE,
1582        ] {
1583            let entry = graph
1584                .iter()
1585                .find(|entry| entry["@id"] == json!(path))
1586                .unwrap_or_else(|| panic!("RO-Crate metadata is missing file entry {path}"));
1587            assert_eq!(entry["sha256"].as_str().map(str::len), Some(64));
1588            assert_eq!(entry["dagml:sha256"].as_str(), entry["sha256"].as_str());
1589            assert_eq!(entry["encodingFormat"], json!("application/json"));
1590            assert!(entry["contentSize"].as_u64().unwrap() > 0);
1591        }
1592        let root = graph
1593            .iter()
1594            .find(|entry| entry["@id"] == json!("./"))
1595            .expect("RO-Crate root dataset is present");
1596        let has_part = root["hasPart"].as_array().unwrap();
1597        assert!(has_part
1598            .iter()
1599            .any(|entry| entry["@id"] == json!(LINEAGE_RECORDS_FILE)));
1600    }
1601
1602    #[test]
1603    fn research_provenance_package_validation_reopens_exported_contracts() {
1604        let plan = fixture_plan();
1605        let bundle = fixture_bundle();
1606        let lineage = vec![fixture_lineage(&plan)];
1607        let package = build_research_provenance_package(
1608            &plan,
1609            &bundle,
1610            &lineage,
1611            &BTreeMap::new(),
1612            None,
1613            None,
1614        )
1615        .unwrap();
1616        let files = package
1617            .files
1618            .iter()
1619            .map(|(path, file)| (path.clone(), file.bytes.clone()))
1620            .collect::<BTreeMap<_, _>>();
1621
1622        let validation = validate_research_provenance_package_files(&files).unwrap();
1623
1624        assert_eq!(
1625            validation.schema_version,
1626            RESEARCH_PROVENANCE_SCHEMA_VERSION
1627        );
1628        assert_eq!(validation.plan_id, plan.id.to_string());
1629        assert_eq!(validation.bundle_id, bundle.bundle_id.to_string());
1630        assert_eq!(validation.file_count, package.files.len());
1631        assert_eq!(validation.checksummed_file_count, package.files.len() - 1);
1632        assert_eq!(validation.lineage_record_count, 1);
1633    }
1634
1635    #[test]
1636    fn research_provenance_package_validation_refuses_tampered_file() {
1637        let plan = fixture_plan();
1638        let bundle = fixture_bundle();
1639        let package =
1640            build_research_provenance_package(&plan, &bundle, &[], &BTreeMap::new(), None, None)
1641                .unwrap();
1642        let mut files = package
1643            .files
1644            .iter()
1645            .map(|(path, file)| (path.clone(), file.bytes.clone()))
1646            .collect::<BTreeMap<_, _>>();
1647        files.insert(LINEAGE_RECORDS_FILE.to_string(), b"[]\n\n".to_vec());
1648
1649        let error = validate_research_provenance_package_files(&files)
1650            .unwrap_err()
1651            .to_string();
1652
1653        assert!(
1654            error.contains("sha256 mismatch"),
1655            "unexpected error: {error}"
1656        );
1657    }
1658
1659    #[test]
1660    fn openlineage_export_uses_validated_research_package() {
1661        let plan = fixture_plan();
1662        let bundle = fixture_bundle();
1663        let lineage = vec![fixture_lineage(&plan)];
1664        let package = build_research_provenance_package(
1665            &plan,
1666            &bundle,
1667            &lineage,
1668            &BTreeMap::new(),
1669            None,
1670            None,
1671        )
1672        .unwrap();
1673        let files = package
1674            .files
1675            .iter()
1676            .map(|(path, file)| (path.clone(), file.bytes.clone()))
1677            .collect::<BTreeMap<_, _>>();
1678
1679        let event = build_openlineage_run_event_from_package_files(
1680            &files,
1681            "dag-ml-test",
1682            "2026-05-27T00:00:00Z",
1683        )
1684        .unwrap();
1685
1686        assert_eq!(event["eventType"], json!("COMPLETE"));
1687        assert_eq!(event["schemaURL"], json!(OPENLINEAGE_RUN_EVENT_SCHEMA_URL));
1688        assert_eq!(event["job"]["namespace"], json!("dag-ml-test"));
1689        assert_eq!(event["run"]["runId"].as_str().map(str::len), Some(36));
1690        assert!(
1691            event["run"]["facets"]["dagml_reproducibility"]["graph_fingerprint"]
1692                .as_str()
1693                .is_some()
1694        );
1695        assert!(!event["inputs"].as_array().unwrap().is_empty());
1696        assert!(event["outputs"]
1697            .as_array()
1698            .unwrap()
1699            .iter()
1700            .any(|output| output["namespace"] == json!("dagml:bundle")));
1701    }
1702
1703    #[test]
1704    fn research_provenance_export_refuses_unknown_lineage_node() {
1705        let plan = fixture_plan();
1706        let bundle = fixture_bundle();
1707        let mut lineage = fixture_lineage(&plan);
1708        lineage.node_id = NodeId::new("model:missing").unwrap();
1709
1710        let error = build_research_provenance_export(
1711            &plan,
1712            &bundle,
1713            &[lineage],
1714            &BTreeMap::new(),
1715            None,
1716            None,
1717        )
1718        .unwrap_err()
1719        .to_string();
1720
1721        assert!(error.contains("unknown node"), "unexpected error: {error}");
1722    }
1723
1724    #[test]
1725    fn research_provenance_export_refuses_mismatched_artifact_manifest() {
1726        let plan = fixture_plan();
1727        let bundle = fixture_bundle();
1728        let mut manifest = FileArtifactManifest {
1729            bundle_id: bundle.bundle_id.clone(),
1730            schema_version: crate::runtime::FILE_ARTIFACT_MANIFEST_SCHEMA_VERSION,
1731            artifacts: Vec::new(),
1732        };
1733        manifest.bundle_id = crate::ids::BundleId::new("bundle:wrong").unwrap();
1734
1735        let error = build_research_provenance_export(
1736            &plan,
1737            &bundle,
1738            &[],
1739            &BTreeMap::new(),
1740            None,
1741            Some(&manifest),
1742        )
1743        .unwrap_err()
1744        .to_string();
1745
1746        assert!(
1747            error.contains("does not match bundle"),
1748            "unexpected error: {error}"
1749        );
1750    }
1751
1752    #[test]
1753    fn research_provenance_export_refuses_unknown_lineage_controller() {
1754        let plan = fixture_plan();
1755        let bundle = fixture_bundle();
1756        let mut lineage = fixture_lineage(&plan);
1757        lineage.controller_id = ControllerId::new("controller:missing").unwrap();
1758
1759        let error = build_research_provenance_export(
1760            &plan,
1761            &bundle,
1762            &[lineage],
1763            &BTreeMap::new(),
1764            None,
1765            None,
1766        )
1767        .unwrap_err()
1768        .to_string();
1769
1770        assert!(
1771            error.contains("unknown controller"),
1772            "unexpected error: {error}"
1773        );
1774    }
1775}